From 03184d98f4d31f431ff2ace2ac1f11a18e43c285 Mon Sep 17 00:00:00 2001 From: "jameskinyua590@gmail.com" <20414083+JayKay24@users.noreply.github.com> Date: Thu, 18 Jun 2026 15:48:05 +0300 Subject: [PATCH 1/6] feat: implement Kafka-to-Spark ingestion pipeline with Docker and producer utilities Signed-off-by: jameskinyua590@gmail.com <20414083+JayKay24@users.noreply.github.com> --- .gitignore | 1 + 3rdparty/requirements.txt | 2 + 3rdparty/user_reqs.lock | 362 +++++++++++++++++- README.md | 2 + agents.md | 7 + projects/ingestion/BUILD | 3 + projects/ingestion/config/input_config.yml | 8 + projects/ingestion/docker/docker-compose.yml | 27 ++ .../ingestion/input_data/user_events.json | 3 + projects/ingestion/json_producer.py | 25 ++ projects/ingestion/kafka_json_to_file_job.py | 74 ++++ 11 files changed, 512 insertions(+), 2 deletions(-) create mode 100644 projects/ingestion/BUILD create mode 100644 projects/ingestion/config/input_config.yml create mode 100644 projects/ingestion/docker/docker-compose.yml create mode 100644 projects/ingestion/input_data/user_events.json create mode 100644 projects/ingestion/json_producer.py create mode 100644 projects/ingestion/kafka_json_to_file_job.py 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..9e03959 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 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/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/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..5c1f165 --- /dev/null +++ b/projects/ingestion/json_producer.py @@ -0,0 +1,25 @@ +import json +import os +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 open(INPUT_FILE_PATH, "r", encoding="utf-8") as f: + records = [json.loads(line) for line in f] + +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..cf89b1c --- /dev/null +++ b/projects/ingestion/kafka_json_to_file_job.py @@ -0,0 +1,74 @@ +import argparse +import os +import yaml +from pyspark.sql import SparkSession +from pyspark.sql.functions import col, from_json + + +def load_config(path): + with open(path, "r", encoding="utf-8") as f: + return yaml.safe_load(f) + + +def run_ingestion_job(config_path: str, output_dir: str): + config = load_config(config_path) + source_conf = config["data_sources"][0] + + 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() + ) + + kafka_options = { + "kafka.bootstrap.servers": source_conf["kafka_config"]["bootstrap_servers"], + "subscribe": source_conf["kafka_config"]["topic"], + "startingOffsets": source_conf["kafka_config"].get( + "starting_offsets", "earliest" + ), + } + + json_schema_str = source_conf.get("json_schema") + if not json_schema_str: + raise ValueError("Missing 'json_schema' in config") + + 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.*") + + df_parsed.write.mode("overwrite").json(output_dir) + + 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.", + ) + + args = parser.parse_args() + run_ingestion_job(args.config, args.output) From 3379e49e029946aca162d7745887501711028da4 Mon Sep 17 00:00:00 2001 From: "jameskinyua590@gmail.com" <20414083+JayKay24@users.noreply.github.com> Date: Thu, 18 Jun 2026 16:16:30 +0300 Subject: [PATCH 2/6] docs: add notebook cell markers and docstrings to Kafka ingestion script Signed-off-by: jameskinyua590@gmail.com <20414083+JayKay24@users.noreply.github.com> --- projects/ingestion/kafka_json_to_file_job.py | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/projects/ingestion/kafka_json_to_file_job.py b/projects/ingestion/kafka_json_to_file_job.py index cf89b1c..cd18247 100644 --- a/projects/ingestion/kafka_json_to_file_job.py +++ b/projects/ingestion/kafka_json_to_file_job.py @@ -1,3 +1,8 @@ +# %% [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 @@ -5,12 +10,23 @@ from pyspark.sql.functions import col, from_json +# %% def load_config(path): + """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): + """Runs Spark Streaming 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. + """ config = load_config(config_path) source_conf = config["data_sources"][0] @@ -49,6 +65,7 @@ def run_ingestion_job(config_path: str, output_dir: str): 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") From 2f47b0edde27cc9adf35da84a3677e2600bf81b5 Mon Sep 17 00:00:00 2001 From: "jameskinyua590@gmail.com" <20414083+JayKay24@users.noreply.github.com> Date: Thu, 18 Jun 2026 16:50:47 +0300 Subject: [PATCH 3/6] docs: add project-level documentation for essentials and ingestion modules in the main repository README Signed-off-by: jameskinyua590@gmail.com <20414083+JayKay24@users.noreply.github.com> --- README.md | 9 ++++++++ projects/essentials/README.md | 34 +++++++++++++++++++++++++++ projects/ingestion/README.md | 43 +++++++++++++++++++++++++++++++++++ 3 files changed, 86 insertions(+) create mode 100644 projects/essentials/README.md create mode 100644 projects/ingestion/README.md diff --git a/README.md b/README.md index 9e03959..92262a2 100644 --- a/README.md +++ b/README.md @@ -44,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/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/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/`. From c004802c8d118bdd1889bf6280dd09c8364496c3 Mon Sep 17 00:00:00 2001 From: "jameskinyua590@gmail.com" <20414083+JayKay24@users.noreply.github.com> Date: Thu, 18 Jun 2026 17:42:01 +0300 Subject: [PATCH 4/6] feat: add descriptive comments to Kafka-to-file ingestion job process Signed-off-by: jameskinyua590@gmail.com <20414083+JayKay24@users.noreply.github.com> --- projects/ingestion/kafka_json_to_file_job.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/projects/ingestion/kafka_json_to_file_job.py b/projects/ingestion/kafka_json_to_file_job.py index cd18247..7ddb884 100644 --- a/projects/ingestion/kafka_json_to_file_job.py +++ b/projects/ingestion/kafka_json_to_file_job.py @@ -27,9 +27,11 @@ def run_ingestion_job(config_path: str, output_dir: str): config_path (str): Path to the ingestion YAML configuration file. output_dir (str): Target directory to save the output JSON events. """ + # Load Kafka connection parameters and schema from config config = load_config(config_path) source_conf = config["data_sources"][0] + # Initialize PySpark session with Kafka integration spark = ( SparkSession.builder.appName("KafkaJsonToFile") .master("local[*]") @@ -40,6 +42,7 @@ def run_ingestion_job(config_path: str, output_dir: str): .getOrCreate() ) + # Set up Kafka connection options kafka_options = { "kafka.bootstrap.servers": source_conf["kafka_config"]["bootstrap_servers"], "subscribe": source_conf["kafka_config"]["topic"], @@ -52,14 +55,18 @@ def run_ingestion_job(config_path: str, output_dir: str): if not json_schema_str: raise ValueError("Missing 'json_schema' in config") + # Read raw binary data from Kafka topic into a DataFrame df_raw = spark.read.format("kafka").options(**kafka_options).load() + # Cast raw Kafka message payload from bytes to a JSON string df_json = df_raw.selectExpr("CAST(value AS STRING) as json_str") + # Parse JSON strings into structured columns based on schema df_parsed = df_json.select( from_json(col("json_str"), json_schema_str).alias("data") ).select("data.*") + # Write the structured data to the local filesystem df_parsed.write.mode("overwrite").json(output_dir) spark.stop() From 19e593f39e497ca6c2799c4bde04f9504ca9d293 Mon Sep 17 00:00:00 2001 From: "jameskinyua590@gmail.com" <20414083+JayKay24@users.noreply.github.com> Date: Thu, 18 Jun 2026 18:09:58 +0300 Subject: [PATCH 5/6] feat: add robust error handling to producer and support for both batch and streaming modes in Spark job Signed-off-by: jameskinyua590@gmail.com <20414083+JayKay24@users.noreply.github.com> --- projects/ingestion/json_producer.py | 23 +++- projects/ingestion/kafka_json_to_file_job.py | 105 +++++++++++++------ 2 files changed, 91 insertions(+), 37 deletions(-) diff --git a/projects/ingestion/json_producer.py b/projects/ingestion/json_producer.py index 5c1f165..324c98a 100644 --- a/projects/ingestion/json_producer.py +++ b/projects/ingestion/json_producer.py @@ -1,5 +1,6 @@ import json import os +import sys from confluent_kafka import Producer # Resolve paths relative to this script @@ -10,9 +11,25 @@ producer_conf = {"bootstrap.servers": "localhost:9092"} producer = Producer(producer_conf) -# Read events from input file -with open(INPUT_FILE_PATH, "r", encoding="utf-8") as f: - records = [json.loads(line) for line in f] +# 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" diff --git a/projects/ingestion/kafka_json_to_file_job.py b/projects/ingestion/kafka_json_to_file_job.py index 7ddb884..f91e11f 100644 --- a/projects/ingestion/kafka_json_to_file_job.py +++ b/projects/ingestion/kafka_json_to_file_job.py @@ -18,20 +18,34 @@ def load_config(path): # %% -def run_ingestion_job(config_path: str, output_dir: str): - """Runs Spark Streaming ingestion job to extract JSON events from Kafka, +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. """ - # Load Kafka connection parameters and schema from config config = load_config(config_path) - source_conf = config["data_sources"][0] - # Initialize PySpark session with Kafka integration + # 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[*]") @@ -42,34 +56,52 @@ def run_ingestion_job(config_path: str, output_dir: str): .getOrCreate() ) - # Set up Kafka connection options - kafka_options = { - "kafka.bootstrap.servers": source_conf["kafka_config"]["bootstrap_servers"], - "subscribe": source_conf["kafka_config"]["topic"], - "startingOffsets": source_conf["kafka_config"].get( - "starting_offsets", "earliest" - ), - } - - json_schema_str = source_conf.get("json_schema") - if not json_schema_str: - raise ValueError("Missing 'json_schema' in config") - - # Read raw binary data from Kafka topic into a DataFrame - df_raw = spark.read.format("kafka").options(**kafka_options).load() - - # Cast raw Kafka message payload from bytes to a JSON string - df_json = df_raw.selectExpr("CAST(value AS STRING) as json_str") - - # Parse JSON strings into structured columns based on schema - df_parsed = df_json.select( - from_json(col("json_str"), json_schema_str).alias("data") - ).select("data.*") - - # Write the structured data to the local filesystem - df_parsed.write.mode("overwrite").json(output_dir) - - spark.stop() + 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: + query.awaitTermination() + except KeyboardInterrupt: + print("Streaming query interrupted by user. Stopping...") + 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 control partition count for small datasets + df_parsed.coalesce(1).write.mode("overwrite").json(output_dir) + finally: + spark.stop() # %% @@ -93,6 +125,11 @@ def run_ingestion_job(config_path: str, output_dir: 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) + run_ingestion_job(args.config, args.output, is_streaming=not args.batch) From a414e3e33705ed54e646a6830d1108e481cb1132 Mon Sep 17 00:00:00 2001 From: "jameskinyua590@gmail.com" <20414083+JayKay24@users.noreply.github.com> Date: Thu, 18 Jun 2026 18:23:32 +0300 Subject: [PATCH 6/6] refactor: add type hints to load_config and improve interruption handling and partitioning comments in Spark job Signed-off-by: jameskinyua590@gmail.com <20414083+JayKay24@users.noreply.github.com> --- projects/ingestion/kafka_json_to_file_job.py | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) diff --git a/projects/ingestion/kafka_json_to_file_job.py b/projects/ingestion/kafka_json_to_file_job.py index f91e11f..1a96c85 100644 --- a/projects/ingestion/kafka_json_to_file_job.py +++ b/projects/ingestion/kafka_json_to_file_job.py @@ -11,7 +11,7 @@ # %% -def load_config(path): +def load_config(path: str) -> dict: """Loads YAML configuration file.""" with open(path, "r", encoding="utf-8") as f: return yaml.safe_load(f) @@ -87,9 +87,16 @@ def run_ingestion_job(config_path: str, output_dir: str, is_streaming: bool = Tr .start() ) try: + # Blocks the main thread, keeping the streaming query active until interrupted query.awaitTermination() except KeyboardInterrupt: - print("Streaming query interrupted by user. Stopping...") + # 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() @@ -98,7 +105,9 @@ def run_ingestion_job(config_path: str, output_dir: str, is_streaming: bool = Tr from_json(col("json_str"), json_schema_str).alias("data") ).select("data.*") - # Coalesce to control partition count for small datasets + # 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()