diff --git a/.gitignore b/.gitignore index 38b5a7c..573e513 100644 --- a/.gitignore +++ b/.gitignore @@ -21,6 +21,7 @@ pants.requirements.pex .DS_Store output_data +output_json # Jupyter Notebooks .ipynb_checkpoints/ diff --git a/3rdparty/requirements.txt b/3rdparty/requirements.txt index 2dd8216..4b2b4fc 100644 --- a/3rdparty/requirements.txt +++ b/3rdparty/requirements.txt @@ -1,3 +1,5 @@ # Third-party Python dependencies for the data engineering projects. pyspark>=3.3.0,<4.0.0 pre-commit>=3.0.0 +pyyaml +confluent-kafka[jsonschema] diff --git a/3rdparty/user_reqs.lock b/3rdparty/user_reqs.lock index b6c7bda..9064edd 100644 --- a/3rdparty/user_reqs.lock +++ b/3rdparty/user_reqs.lock @@ -9,8 +9,10 @@ // "CPython==3.10.9" // ], // "generated_with_requirements": [ +// "confluent-kafka[jsonschema]", // "pre-commit>=3.0.0", -// "pyspark<4.0.0,>=3.3.0" +// "pyspark<4.0.0,>=3.3.0", +// "pyyaml" // ], // "manylinux": "manylinux2014", // "requirement_constraints": [], @@ -46,6 +48,360 @@ "requires_python": ">=3.10", "version": "3.5.0" }, + { + "artifacts": [ + { + "algorithm": "sha256", + "hash": "4b9d6e0170aed7856553429d1d2b3e9373a4c4ba0c264204e05a947a52cf2c2f", + "url": "https://files.pythonhosted.org/packages/f8/b1/37f112f18e7cee6bc90eaadb475e36f80fabb1a9d3d7acbe0310e1c5326e/confluent_kafka-2.14.2-cp310-cp310-manylinux_2_28_x86_64.whl" + }, + { + "algorithm": "sha256", + "hash": "9f6f9718c0cd8d32a61aa58010c34b741dd62ec2ed9f470fcff36f7ac3d1006e", + "url": "https://files.pythonhosted.org/packages/14/64/13d093f6fc6b318c301942f579599cc1cae07a93061442fb1a81e40962e4/confluent_kafka-2.14.2-cp310-cp310-macosx_10_9_x86_64.whl" + }, + { + "algorithm": "sha256", + "hash": "5714f9402ef7e07c42bb1d3e175c9dd811a6b5877b522828630ae800503deaf0", + "url": "https://files.pythonhosted.org/packages/96/9f/44f326e5073b618712b4b52575d72c093da896642233080c208e2ab0025b/confluent_kafka-2.14.2-cp310-cp310-macosx_11_0_arm64.whl" + }, + { + "algorithm": "sha256", + "hash": "4c9b1c4fe8f4364a5828dc35fb481577a7c0b6960876f5e53a0370705ca0247d", + "url": "https://files.pythonhosted.org/packages/98/a8/4f3136f8f87eda871003aecac8c57f4aabfd9cb600f55699a32c2fb225d8/confluent_kafka-2.14.2-cp310-cp310-manylinux_2_28_aarch64.whl" + }, + { + "algorithm": "sha256", + "hash": "fc827265571a778b1ff560ab2f3ec5dce3c573c29eb4cf6ded9718d55a5e4262", + "url": "https://files.pythonhosted.org/packages/ff/e5/58ac277094b03a51d9725798d0e53a60c1be69dc4330ae9cedc2d9efa430/confluent_kafka-2.14.2.tar.gz" + } + ], + "project_name": "confluent-kafka", + "requires_dists": [ + "async-timeout; extra == \"all\"", + "async-timeout; extra == \"dev\"", + "async-timeout; extra == \"tests\"", + "attrs; extra == \"all\"", + "attrs; extra == \"all\"", + "attrs; extra == \"dev\"", + "attrs; extra == \"dev\"", + "attrs; extra == \"examples\"", + "attrs; extra == \"tests\"", + "attrs>=21.2.0; extra == \"all\"", + "attrs>=21.2.0; extra == \"avro\"", + "attrs>=21.2.0; extra == \"dev\"", + "attrs>=21.2.0; extra == \"docs\"", + "attrs>=21.2.0; extra == \"json\"", + "attrs>=21.2.0; extra == \"protobuf\"", + "attrs>=21.2.0; extra == \"rules\"", + "attrs>=21.2.0; extra == \"schema-registry\"", + "attrs>=21.2.0; extra == \"schemaregistry\"", + "attrs>=21.2.0; extra == \"tests\"", + "authlib>=1.0.0; extra == \"all\"", + "authlib>=1.0.0; extra == \"all\"", + "authlib>=1.0.0; extra == \"avro\"", + "authlib>=1.0.0; extra == \"dev\"", + "authlib>=1.0.0; extra == \"dev\"", + "authlib>=1.0.0; extra == \"docs\"", + "authlib>=1.0.0; extra == \"examples\"", + "authlib>=1.0.0; extra == \"json\"", + "authlib>=1.0.0; extra == \"protobuf\"", + "authlib>=1.0.0; extra == \"rules\"", + "authlib>=1.0.0; extra == \"schema-registry\"", + "authlib>=1.0.0; extra == \"schemaregistry\"", + "authlib>=1.0.0; extra == \"tests\"", + "avro<2,>=1.11.1; extra == \"all\"", + "avro<2,>=1.11.1; extra == \"all\"", + "avro<2,>=1.11.1; extra == \"avro\"", + "avro<2,>=1.11.1; extra == \"dev\"", + "avro<2,>=1.11.1; extra == \"dev\"", + "avro<2,>=1.11.1; extra == \"docs\"", + "avro<2,>=1.11.1; extra == \"examples\"", + "avro<2,>=1.11.1; extra == \"tests\"", + "azure-identity; extra == \"all\"", + "azure-identity; extra == \"all\"", + "azure-identity; extra == \"dev\"", + "azure-identity; extra == \"dev\"", + "azure-identity; extra == \"docs\"", + "azure-identity; extra == \"examples\"", + "azure-identity; extra == \"rules\"", + "azure-identity; extra == \"tests\"", + "azure-keyvault-keys; extra == \"all\"", + "azure-keyvault-keys; extra == \"all\"", + "azure-keyvault-keys; extra == \"dev\"", + "azure-keyvault-keys; extra == \"dev\"", + "azure-keyvault-keys; extra == \"docs\"", + "azure-keyvault-keys; extra == \"examples\"", + "azure-keyvault-keys; extra == \"rules\"", + "azure-keyvault-keys; extra == \"tests\"", + "black>=24.0.0; extra == \"all\"", + "black>=24.0.0; extra == \"dev\"", + "black>=24.0.0; extra == \"tests\"", + "boto3; extra == \"all\"", + "boto3; extra == \"dev\"", + "boto3; extra == \"examples\"", + "boto3>=1.35; extra == \"all\"", + "boto3>=1.35; extra == \"dev\"", + "boto3>=1.35; extra == \"docs\"", + "boto3>=1.35; extra == \"rules\"", + "boto3>=1.35; extra == \"tests\"", + "cachetools; extra == \"all\"", + "cachetools; extra == \"dev\"", + "cachetools; extra == \"examples\"", + "cachetools>=5.5.0; extra == \"all\"", + "cachetools>=5.5.0; extra == \"avro\"", + "cachetools>=5.5.0; extra == \"dev\"", + "cachetools>=5.5.0; extra == \"docs\"", + "cachetools>=5.5.0; extra == \"json\"", + "cachetools>=5.5.0; extra == \"protobuf\"", + "cachetools>=5.5.0; extra == \"rules\"", + "cachetools>=5.5.0; extra == \"schema-registry\"", + "cachetools>=5.5.0; extra == \"schemaregistry\"", + "cachetools>=5.5.0; extra == \"tests\"", + "cel-python>=0.4.0; extra == \"all\"", + "cel-python>=0.4.0; extra == \"all\"", + "cel-python>=0.4.0; extra == \"dev\"", + "cel-python>=0.4.0; extra == \"dev\"", + "cel-python>=0.4.0; extra == \"docs\"", + "cel-python>=0.4.0; extra == \"examples\"", + "cel-python>=0.4.0; extra == \"rules\"", + "cel-python>=0.4.0; extra == \"tests\"", + "certifi; extra == \"all\"", + "certifi; extra == \"avro\"", + "certifi; extra == \"dev\"", + "certifi; extra == \"docs\"", + "certifi; extra == \"json\"", + "certifi; extra == \"protobuf\"", + "certifi; extra == \"rules\"", + "certifi; extra == \"schema-registry\"", + "certifi; extra == \"schemaregistry\"", + "certifi; extra == \"tests\"", + "confluent-kafka; extra == \"all\"", + "confluent-kafka; extra == \"dev\"", + "confluent-kafka; extra == \"examples\"", + "fastapi; extra == \"all\"", + "fastapi; extra == \"dev\"", + "fastapi; extra == \"examples\"", + "fastavro<1.8.0,>=1.5.4; python_version == \"3.7\" and extra == \"all\"", + "fastavro<1.8.0,>=1.5.4; python_version == \"3.7\" and extra == \"all\"", + "fastavro<1.8.0,>=1.5.4; python_version == \"3.7\" and extra == \"avro\"", + "fastavro<1.8.0,>=1.5.4; python_version == \"3.7\" and extra == \"dev\"", + "fastavro<1.8.0,>=1.5.4; python_version == \"3.7\" and extra == \"dev\"", + "fastavro<1.8.0,>=1.5.4; python_version == \"3.7\" and extra == \"docs\"", + "fastavro<1.8.0,>=1.5.4; python_version == \"3.7\" and extra == \"examples\"", + "fastavro<1.8.0,>=1.5.4; python_version == \"3.7\" and extra == \"tests\"", + "fastavro<2,>=1.5.4; python_version > \"3.7\" and extra == \"all\"", + "fastavro<2,>=1.5.4; python_version > \"3.7\" and extra == \"all\"", + "fastavro<2,>=1.5.4; python_version > \"3.7\" and extra == \"avro\"", + "fastavro<2,>=1.5.4; python_version > \"3.7\" and extra == \"dev\"", + "fastavro<2,>=1.5.4; python_version > \"3.7\" and extra == \"dev\"", + "fastavro<2,>=1.5.4; python_version > \"3.7\" and extra == \"docs\"", + "fastavro<2,>=1.5.4; python_version > \"3.7\" and extra == \"examples\"", + "fastavro<2,>=1.5.4; python_version > \"3.7\" and extra == \"tests\"", + "flake8; extra == \"all\"", + "flake8; extra == \"dev\"", + "flake8; extra == \"tests\"", + "google-api-core; extra == \"all\"", + "google-api-core; extra == \"all\"", + "google-api-core; extra == \"dev\"", + "google-api-core; extra == \"dev\"", + "google-api-core; extra == \"docs\"", + "google-api-core; extra == \"examples\"", + "google-api-core; extra == \"rules\"", + "google-api-core; extra == \"tests\"", + "google-auth; extra == \"all\"", + "google-auth; extra == \"all\"", + "google-auth; extra == \"dev\"", + "google-auth; extra == \"dev\"", + "google-auth; extra == \"docs\"", + "google-auth; extra == \"examples\"", + "google-auth; extra == \"rules\"", + "google-auth; extra == \"tests\"", + "google-cloud-kms; extra == \"all\"", + "google-cloud-kms; extra == \"all\"", + "google-cloud-kms; extra == \"dev\"", + "google-cloud-kms; extra == \"dev\"", + "google-cloud-kms; extra == \"docs\"", + "google-cloud-kms; extra == \"examples\"", + "google-cloud-kms; extra == \"rules\"", + "google-cloud-kms; extra == \"tests\"", + "google-re2<1.1.20251105; extra == \"all\"", + "google-re2<1.1.20251105; extra == \"dev\"", + "google-re2<1.1.20251105; extra == \"docs\"", + "google-re2<1.1.20251105; extra == \"rules\"", + "google-re2<1.1.20251105; extra == \"tests\"", + "googleapis-common-protos; extra == \"all\"", + "googleapis-common-protos; extra == \"all\"", + "googleapis-common-protos; extra == \"dev\"", + "googleapis-common-protos; extra == \"dev\"", + "googleapis-common-protos; extra == \"docs\"", + "googleapis-common-protos; extra == \"examples\"", + "googleapis-common-protos; extra == \"protobuf\"", + "googleapis-common-protos; extra == \"tests\"", + "hkdf==0.0.3; extra == \"all\"", + "hkdf==0.0.3; extra == \"all\"", + "hkdf==0.0.3; extra == \"dev\"", + "hkdf==0.0.3; extra == \"dev\"", + "hkdf==0.0.3; extra == \"docs\"", + "hkdf==0.0.3; extra == \"examples\"", + "hkdf==0.0.3; extra == \"rules\"", + "hkdf==0.0.3; extra == \"tests\"", + "httpx>=0.26; extra == \"all\"", + "httpx>=0.26; extra == \"all\"", + "httpx>=0.26; extra == \"avro\"", + "httpx>=0.26; extra == \"dev\"", + "httpx>=0.26; extra == \"dev\"", + "httpx>=0.26; extra == \"docs\"", + "httpx>=0.26; extra == \"examples\"", + "httpx>=0.26; extra == \"json\"", + "httpx>=0.26; extra == \"protobuf\"", + "httpx>=0.26; extra == \"rules\"", + "httpx>=0.26; extra == \"schema-registry\"", + "httpx>=0.26; extra == \"schemaregistry\"", + "httpx>=0.26; extra == \"tests\"", + "hvac; extra == \"all\"", + "hvac; extra == \"all\"", + "hvac; extra == \"dev\"", + "hvac; extra == \"dev\"", + "hvac; extra == \"docs\"", + "hvac; extra == \"examples\"", + "hvac; extra == \"rules\"", + "hvac; extra == \"tests\"", + "isort>=5.13.0; extra == \"all\"", + "isort>=5.13.0; extra == \"dev\"", + "isort>=5.13.0; extra == \"tests\"", + "jsonata-python; extra == \"all\"", + "jsonata-python; extra == \"all\"", + "jsonata-python; extra == \"dev\"", + "jsonata-python; extra == \"dev\"", + "jsonata-python; extra == \"docs\"", + "jsonata-python; extra == \"examples\"", + "jsonata-python; extra == \"rules\"", + "jsonata-python; extra == \"tests\"", + "jsonschema>=4.18.0; extra == \"all\"", + "jsonschema>=4.18.0; extra == \"all\"", + "jsonschema>=4.18.0; extra == \"dev\"", + "jsonschema>=4.18.0; extra == \"dev\"", + "jsonschema>=4.18.0; extra == \"docs\"", + "jsonschema>=4.18.0; extra == \"examples\"", + "jsonschema>=4.18.0; extra == \"json\"", + "jsonschema>=4.18.0; extra == \"tests\"", + "mypy; extra == \"all\"", + "mypy; extra == \"dev\"", + "mypy; extra == \"tests\"", + "opentelemetry-distro; extra == \"all\"", + "opentelemetry-distro; extra == \"soaktest\"", + "opentelemetry-exporter-otlp; extra == \"all\"", + "opentelemetry-exporter-otlp; extra == \"soaktest\"", + "orjson; extra == \"all\"", + "orjson; extra == \"dev\"", + "orjson; extra == \"tests\"", + "orjson>=3.10; extra == \"all\"", + "orjson>=3.10; extra == \"all\"", + "orjson>=3.10; extra == \"dev\"", + "orjson>=3.10; extra == \"dev\"", + "orjson>=3.10; extra == \"docs\"", + "orjson>=3.10; extra == \"examples\"", + "orjson>=3.10; extra == \"json\"", + "orjson>=3.10; extra == \"tests\"", + "pandoc; extra == \"all\"", + "pandoc; extra == \"dev\"", + "pandoc; extra == \"docs\"", + "pluggy<1.6.0; extra == \"all\"", + "pluggy<1.6.0; extra == \"dev\"", + "pluggy<1.6.0; extra == \"tests\"", + "protobuf; extra == \"all\"", + "protobuf; extra == \"all\"", + "protobuf; extra == \"dev\"", + "protobuf; extra == \"dev\"", + "protobuf; extra == \"docs\"", + "protobuf; extra == \"examples\"", + "protobuf; extra == \"protobuf\"", + "protobuf; extra == \"tests\"", + "psutil; extra == \"all\"", + "psutil; extra == \"soaktest\"", + "pydantic; extra == \"all\"", + "pydantic; extra == \"dev\"", + "pydantic; extra == \"examples\"", + "pyrsistent; extra == \"all\"", + "pyrsistent; extra == \"all\"", + "pyrsistent; extra == \"dev\"", + "pyrsistent; extra == \"dev\"", + "pyrsistent; extra == \"docs\"", + "pyrsistent; extra == \"examples\"", + "pyrsistent; extra == \"json\"", + "pyrsistent; extra == \"tests\"", + "pytest-asyncio; extra == \"all\"", + "pytest-asyncio; extra == \"dev\"", + "pytest-asyncio; extra == \"tests\"", + "pytest-timeout; extra == \"all\"", + "pytest-timeout; extra == \"dev\"", + "pytest-timeout; extra == \"tests\"", + "pytest; extra == \"all\"", + "pytest; extra == \"dev\"", + "pytest; extra == \"tests\"", + "pytest_cov; extra == \"all\"", + "pytest_cov; extra == \"dev\"", + "pytest_cov; extra == \"tests\"", + "pyyaml>=6.0.0; extra == \"all\"", + "pyyaml>=6.0.0; extra == \"all\"", + "pyyaml>=6.0.0; extra == \"dev\"", + "pyyaml>=6.0.0; extra == \"dev\"", + "pyyaml>=6.0.0; extra == \"docs\"", + "pyyaml>=6.0.0; extra == \"examples\"", + "pyyaml>=6.0.0; extra == \"rules\"", + "pyyaml>=6.0.0; extra == \"tests\"", + "requests-mock; extra == \"all\"", + "requests-mock; extra == \"dev\"", + "requests-mock; extra == \"tests\"", + "requests; extra == \"all\"", + "requests; extra == \"all\"", + "requests; extra == \"avro\"", + "requests; extra == \"dev\"", + "requests; extra == \"dev\"", + "requests; extra == \"docs\"", + "requests; extra == \"examples\"", + "requests; extra == \"tests\"", + "respx; extra == \"all\"", + "respx; extra == \"dev\"", + "respx; extra == \"tests\"", + "six; extra == \"all\"", + "six; extra == \"dev\"", + "six; extra == \"examples\"", + "sphinx-rtd-theme; extra == \"all\"", + "sphinx-rtd-theme; extra == \"dev\"", + "sphinx-rtd-theme; extra == \"docs\"", + "sphinx; extra == \"all\"", + "sphinx; extra == \"dev\"", + "sphinx; extra == \"docs\"", + "tink; extra == \"all\"", + "tink; extra == \"all\"", + "tink; extra == \"dev\"", + "tink; extra == \"dev\"", + "tink; extra == \"docs\"", + "tink; extra == \"examples\"", + "tink; extra == \"rules\"", + "tink; extra == \"tests\"", + "tomli; python_version < \"3.11\" and extra == \"all\"", + "tomli; python_version < \"3.11\" and extra == \"dev\"", + "tomli; python_version < \"3.11\" and extra == \"docs\"", + "types-cachetools; extra == \"all\"", + "types-cachetools; extra == \"dev\"", + "types-cachetools; extra == \"tests\"", + "types-requests; extra == \"all\"", + "types-requests; extra == \"dev\"", + "types-requests; extra == \"tests\"", + "typing-extensions; python_version < \"3.11\"", + "urllib3<3; extra == \"all\"", + "urllib3<3; extra == \"dev\"", + "urllib3<3; extra == \"tests\"", + "uvicorn; extra == \"all\"", + "uvicorn; extra == \"dev\"", + "uvicorn; extra == \"examples\"" + ], + "requires_python": ">=3.8", + "version": "2.14.2" + }, { "artifacts": [ { @@ -343,8 +699,10 @@ "pip_version": "24.0", "prefer_older_binary": false, "requirements": [ + "confluent-kafka[jsonschema]", "pre-commit>=3.0.0", - "pyspark<4.0.0,>=3.3.0" + "pyspark<4.0.0,>=3.3.0", + "pyyaml" ], "requires_python": [ "==3.10.9" diff --git a/README.md b/README.md index 7ef3de4..92262a2 100644 --- a/README.md +++ b/README.md @@ -35,6 +35,8 @@ data-engineering/ │ ├── requirements.txt # Lists project requirements (pandas, PySpark, etc.) │ └── user_reqs.lock # Pants generated dependency lockfile ├── projects/ # Directory containing all sub-projects +│ ├── ingestion/ # Chapter 4 Kafka/Spark ingestion project +│ └── essentials/ # Chapter 2 basic Spark examples └── scripts/ ├── BUILD # Configures scripts targets for Pants └── ai_pr_reviewer.py # Python script that runs Gemini AI code reviews @@ -42,6 +44,15 @@ data-engineering/ --- +## 📁 Projects + +Each project under the `projects/` directory represents a separate learning milestone with self-contained instructions, docker components, and code: + +* [projects/essentials/](projects/essentials/) — Basic local PySpark processing examples (see [projects/essentials/README.md](projects/essentials/README.md)). +* [projects/ingestion/](projects/ingestion/) — Real-time event ingestion using Kafka and Spark Streaming (see [projects/ingestion/README.md](projects/ingestion/README.md)). + +--- + ## 🚀 Getting Started ### 1. Prerequisites diff --git a/agents.md b/agents.md index d8ab861..b4ac7a2 100644 --- a/agents.md +++ b/agents.md @@ -26,6 +26,13 @@ This repository is a **Python Data Engineering Monorepo** managed by the **Pants * `employee_partition_by_hire_date.py`: Local partitioning Spark script. * `input_data/`: Small CSV/txt sample inputs. * `output_data/`: Automatically generated Spark output targets (ignored by git). +* `projects/ingestion/`: Ingestion project (derived implementation from *Hello Modern Data Pipelines*, Chapter 4). + * `config/input_config.yml`: Spark Ingestion configuration YAML file. + * `docker/docker-compose.yml`: Zookeeper, Kafka, and Schema Registry Compose setup. + * `input_data/user_events.json`: Sample event stream dataset. + * `json_producer.py`: Kafka JSON message producer script. + * `kafka_json_to_file_job.py`: PySpark job ingestion script with Spark SQL Kafka integration. + * `BUILD`: Pants build definition for the ingestion project. * `scripts/`: Python utility scripts. * `ai_pr_reviewer.py`: The AI code reviewer script powered by the Gemini API. * `BUILD`: Pants build definition for the scripts directory. diff --git a/projects/essentials/README.md b/projects/essentials/README.md new file mode 100644 index 0000000..22c38f2 --- /dev/null +++ b/projects/essentials/README.md @@ -0,0 +1,34 @@ +# Essentials Project (Chapter 2) + +This project contains initial Spark processing examples derived from Chapter 2 of *Hello Modern Data Pipelines*. It demonstrates basic local batch data processing using PySpark. + +--- + +## 📁 Project Contents + +* [word_count.py](projects/essentials/word_count.py): Processes unstructured text data to perform a classic word-count calculation. +* [employee_partition_by_hire_date.py](projects/essentials/employee_partition_by_hire_date.py): Demonstrates PySpark DataFrame API usage, reading CSV employee data and partitioning the output by hire date. +* `input_data/`: Contains sample CSV and text input files, such as [employee_data.csv](projects/essentials/input_data/employee_data.csv), for testing the scripts. + +--- + +## 🚀 How to Run + +Before running the scripts, ensure your virtual environment is active: +```bash +source .venv/bin/activate +``` + +### 1. Run WordCount +Run the word count script: +```bash +python projects/essentials/word_count.py +``` +This generates the results in `projects/essentials/output_data/word_count/`. + +### 2. Run Employee Partitioning +Run the partitioning script: +```bash +python projects/essentials/employee_partition_by_hire_date.py +``` +This partitions the employee data and writes it to `projects/essentials/output_data/employee_partition/`. diff --git a/projects/ingestion/BUILD b/projects/ingestion/BUILD new file mode 100644 index 0000000..4a902ca --- /dev/null +++ b/projects/ingestion/BUILD @@ -0,0 +1,3 @@ +python_sources( + name="lib", +) diff --git a/projects/ingestion/README.md b/projects/ingestion/README.md new file mode 100644 index 0000000..dc6f789 --- /dev/null +++ b/projects/ingestion/README.md @@ -0,0 +1,43 @@ +# Ingestion Project (Chapter 4) + +This project implements the advanced data ingestion and integration patterns derived from Chapter 4 of *Hello Modern Data Pipelines*. It demonstrates real-time integration by producing events to a Kafka broker and streaming/reading them into Spark. + +--- + +## 📁 Project Contents + +* [docker/docker-compose.yml](projects/ingestion/docker/docker-compose.yml): Launches local Zookeeper, Kafka, and Schema Registry containers. +* [json_producer.py](projects/ingestion/json_producer.py): Python producer script that publishes mock events to the Kafka broker. +* [kafka_json_to_file_job.py](projects/ingestion/kafka_json_to_file_job.py): PySpark streaming/batch job that pulls events from Kafka, parses them with a schema, and saves them locally as JSON. +* [config/input_config.yml](projects/ingestion/config/input_config.yml): Configuration file specifying broker connections, target topics, and JSON parsing schema. +* [input_data/user_events.json](projects/ingestion/input_data/user_events.json): Sample JSON data file with mock events. + +--- + +## 🚀 How to Run + +Ensure your virtual environment is active: +```bash +source .venv/bin/activate +``` + +### 1. Spin up Kafka Infrastucture +Start the Zookeeper, Kafka, and Schema Registry containers: +```bash +cd projects/ingestion/docker +docker-compose up -d +cd ../../.. +``` + +### 2. Produce mock JSON events +Publish the sample events from `user_events.json` into Kafka: +```bash +python projects/ingestion/json_producer.py +``` + +### 3. Run Ingestion Spark Job +Extract the events from Kafka and write them to output directories: +```bash +python projects/ingestion/kafka_json_to_file_job.py +``` +Outputs are written locally to `projects/ingestion/output_json/user_events/`. diff --git a/projects/ingestion/config/input_config.yml b/projects/ingestion/config/input_config.yml new file mode 100644 index 0000000..f399e01 --- /dev/null +++ b/projects/ingestion/config/input_config.yml @@ -0,0 +1,8 @@ +data_sources: + - source_id: user_events_kafka + source_type: kafka + kafka_config: + bootstrap_servers: "localhost:9092" + topic: "user-events-json" + starting_offsets: "earliest" + json_schema: "STRUCT" diff --git a/projects/ingestion/docker/docker-compose.yml b/projects/ingestion/docker/docker-compose.yml new file mode 100644 index 0000000..db3a144 --- /dev/null +++ b/projects/ingestion/docker/docker-compose.yml @@ -0,0 +1,27 @@ +version: '2' +services: + zookeeper: + image: confluentinc/cp-zookeeper:7.2.1 + environment: + ZOOKEEPER_CLIENT_PORT: 2181 + + kafka: + image: confluentinc/cp-kafka:7.2.1 + ports: + - "9092:9092" + environment: + KAFKA_BROKER_ID: 1 + KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092 + KAFKA_LISTENERS: PLAINTEXT://0.0.0.0:29092,PLAINTEXT_HOST://0.0.0.0:9092 + KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + + schema-registry: + image: confluentinc/cp-schema-registry:7.2.1 + ports: + - "8081:8081" + environment: + SCHEMA_REGISTRY_HOST_NAME: schema-registry + SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://kafka:29092 diff --git a/projects/ingestion/input_data/user_events.json b/projects/ingestion/input_data/user_events.json new file mode 100644 index 0000000..5073316 --- /dev/null +++ b/projects/ingestion/input_data/user_events.json @@ -0,0 +1,3 @@ +{"user_id": 101, "event_type": "login", "timestamp": "2025-07-07T09:00:00Z"} +{"user_id": 102, "event_type": "logout", "timestamp": "2025-07-07T09:05:00Z"} +{"user_id": 104, "event_type": "login", "timestamp": "2025-07-07T09:05:00Z"} diff --git a/projects/ingestion/json_producer.py b/projects/ingestion/json_producer.py new file mode 100644 index 0000000..324c98a --- /dev/null +++ b/projects/ingestion/json_producer.py @@ -0,0 +1,42 @@ +import json +import os +import sys +from confluent_kafka import Producer + +# Resolve paths relative to this script +SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__)) +INPUT_FILE_PATH = os.path.join(SCRIPT_DIR, "input_data/user_events.json") + +# Kafka producer configuration +producer_conf = {"bootstrap.servers": "localhost:9092"} +producer = Producer(producer_conf) + +# Read events from input file with error handling +records = [] +try: + with open(INPUT_FILE_PATH, "r", encoding="utf-8") as f: + for line in f: + if line.strip(): + try: + records.append(json.loads(line)) + except json.JSONDecodeError as e: + print( + f"Skipping malformed JSON line: {line.strip()} - Error: {e}", + file=sys.stderr, + ) +except FileNotFoundError: + print(f"Error: Input file not found at {INPUT_FILE_PATH}", file=sys.stderr) + sys.exit(1) +except IOError as e: + print(f"Error reading input file {INPUT_FILE_PATH}: {e}", file=sys.stderr) + sys.exit(1) + +topic = "user-events-json" + +# Produce each record to Kafka +for record in records: + json_str = json.dumps(record) + producer.produce(topic=topic, value=json_str) + print("Produced:", json_str) + +producer.flush() diff --git a/projects/ingestion/kafka_json_to_file_job.py b/projects/ingestion/kafka_json_to_file_job.py new file mode 100644 index 0000000..1a96c85 --- /dev/null +++ b/projects/ingestion/kafka_json_to_file_job.py @@ -0,0 +1,144 @@ +# %% [markdown] +# # Kafka JSON to File Ingestion Job +# Extracts user events from a Kafka topic, parses them according to a schema config, and writes them to the local filesystem. + +# %% +import argparse +import os +import yaml +from pyspark.sql import SparkSession +from pyspark.sql.functions import col, from_json + + +# %% +def load_config(path: str) -> dict: + """Loads YAML configuration file.""" + with open(path, "r", encoding="utf-8") as f: + return yaml.safe_load(f) + + +# %% +def run_ingestion_job(config_path: str, output_dir: str, is_streaming: bool = True): + """Runs Spark Ingestion job to extract JSON events from Kafka, + + parses them using a JSON schema, and writes them to output_dir. + + Args: + config_path (str): Path to the ingestion YAML configuration file. + output_dir (str): Target directory to save the output JSON events. + is_streaming (bool): If True, run as Structured Streaming job. If False, run as Batch job. + """ + config = load_config(config_path) + + # Configuration Validation + if not isinstance(config, dict) or "data_sources" not in config: + raise ValueError("Invalid configuration file structure: missing 'data_sources'") + + data_sources = config["data_sources"] + if not data_sources or not isinstance(data_sources, list): + raise ValueError("'data_sources' must be a non-empty list") + + source_conf = data_sources[0] + if "kafka_config" not in source_conf: + raise ValueError("Missing 'kafka_config' in configuration source") + + kafka_conf = source_conf["kafka_config"] + if "bootstrap_servers" not in kafka_conf or "topic" not in kafka_conf: + raise ValueError("Missing 'bootstrap_servers' or 'topic' in 'kafka_config'") + + spark = ( + SparkSession.builder.appName("KafkaJsonToFile") + .master("local[*]") + .config( + "spark.jars.packages", + "org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.8", + ) + .getOrCreate() + ) + + try: + kafka_options = { + "kafka.bootstrap.servers": kafka_conf["bootstrap_servers"], + "subscribe": kafka_conf["topic"], + "startingOffsets": kafka_conf.get("starting_offsets", "earliest"), + } + + json_schema_str = source_conf.get("json_schema") + if not json_schema_str: + raise ValueError("Missing 'json_schema' in config") + + if is_streaming: + # Structured Streaming Read + df_raw = spark.readStream.format("kafka").options(**kafka_options).load() + df_json = df_raw.selectExpr("CAST(value AS STRING) as json_str") + df_parsed = df_json.select( + from_json(col("json_str"), json_schema_str).alias("data") + ).select("data.*") + + checkpoint_path = os.path.join(output_dir, "_checkpoint") + + # Structured Streaming Write (using Append mode) + query = ( + df_parsed.writeStream.outputMode("append") + .format("json") + .option("path", output_dir) + .option("checkpointLocation", checkpoint_path) + .trigger(processingTime="10 seconds") + .start() + ) + try: + # Blocks the main thread, keeping the streaming query active until interrupted + query.awaitTermination() + except KeyboardInterrupt: + # Triggered when the user presses Ctrl+C in the terminal + import sys + + print( + "Streaming query interrupted by user. Stopping...", + file=sys.stderr, + ) + else: + # Batch Read + df_raw = spark.read.format("kafka").options(**kafka_options).load() + df_json = df_raw.selectExpr("CAST(value AS STRING) as json_str") + df_parsed = df_json.select( + from_json(col("json_str"), json_schema_str).alias("data") + ).select("data.*") + + # Coalesce to 1 partition for small datasets to avoid many small files. + # For larger datasets, remove .coalesce(1) to let Spark manage partitions + # or repartition based on business keys. + df_parsed.coalesce(1).write.mode("overwrite").json(output_dir) + finally: + spark.stop() + + +# %% +if __name__ == "__main__": + script_dir = os.path.dirname(os.path.abspath(__file__)) + default_config = os.path.join(script_dir, "config/input_config.yml") + default_output = os.path.join(script_dir, "output_json/user_events") + + parser = argparse.ArgumentParser( + description="Extract user events from Kafka and save as JSON." + ) + parser.add_argument( + "--config", + type=str, + default=default_config, + help="Path to the input configuration YAML file.", + ) + parser.add_argument( + "--output", + type=str, + default=default_output, + help="Path to save the output JSON events.", + ) + parser.add_argument( + "--batch", + action="store_true", + help="Run as a one-time batch job instead of a continuous streaming job.", + ) + + args = parser.parse_args() + run_ingestion_job(args.config, args.output, is_streaming=not args.batch)