How to Implement Real-Time Streaming with Kafka and Flink: A Complete PyFlink Tutorial
Implement real-time streaming with Kafka and Flink by configuring a Kafka source table in PyFlink SQL, enriching events with a Python user-defined function (UDF) for geolocation lookup, and sinking the transformed data to PostgreSQL via a continuous INSERT query orchestrated through Docker Compose.
To implement real-time streaming with Kafka and Flink, you need a robust architecture that connects stream sources, processing logic, and persistent storage. The DataExpert-io/data-engineer-handbook repository provides a production-ready example that demonstrates exactly this pattern using PyFlink, Confluent Kafka, and PostgreSQL in a containerized environment. This tutorial walks through the end-to-end implementation, from the Docker Compose infrastructure to the PyFlink job logic that powers the pipeline.
Architecture Overview
The pipeline ingests web traffic events from a Kafka topic, enriches each record with geographic metadata based on the client IP address, and persists the results to a relational database. According to the DataExpert-io/data-engineer-handbook source code, the architecture relies on three core components: a Kafka source table, a Python UDF for external API calls, and a JDBC sink table.
Setting Up the Streaming Environment
Before deploying the Flink job, you must configure the infrastructure that hosts the Flink cluster and PostgreSQL instance.
Docker Compose Configuration
The docker-compose.yml file defines a Flink JobManager, TaskManager, and PostgreSQL service, ensuring all components share a networked environment. The JobManager exposes the Flink Web UI on port 8081 and orchestrates job deployment, while the TaskManager executes the actual stream processing tasks. You can view the complete service definitions in the repository's compose file at [docker-compose.yml](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/docker-compose.yml#L1-L54).
Implementing the Flink Job
The core logic resides in src/job/start_job.py, which initializes the StreamExecutionEnvironment, registers tables, and executes a continuous SQL query.
Configuring the Kafka Source
The create_events_source_kafka function defines a table that connects to a Confluent-hosted Kafka cluster using SASL_SSL authentication. It parses JSON messages from the web-traffic topic and converts the event_time string field into a proper timestamp using the TO_TIMESTAMP function. The implementation specifies the Kafka consumer group, security protocol, and JSON format within the SQL DDL, as shown in lines 83 through 100 of [start_job.py](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py#L83-L100).
Creating the PostgreSQL Sink
To persist enriched data, the create_processed_events_sink_postgres function registers a JDBC sink table named processed_events. This table stores the IP address, event timestamp, referrer, host, URL, and geolocation JSON. The connector configuration includes the PostgreSQL URL, driver, and credentials injected via environment variables. See the sink table definition in lines 36 through 56 of [start_job.py](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py#L36-L56).
Enriching Data with a Python UDF
The GetLocation class extends PyFlink's ScalarFunction and calls the ip2location.io API to resolve IP addresses into country, state, and city metadata. Registered as the UDF get_location, this function executes for each record flowing through the pipeline. The UDF implementation and registration appear in lines 58 through 80 of [start_job.py](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py#L58-L80).
Orchestrating the Continuous Query
The log_processing function ties the components together by executing an INSERT INTO processed_events SELECT ... FROM events statement. This query runs continuously, reading from the Kafka source, applying the get_location UDF to each row, and writing results to the PostgreSQL sink. The job configures checkpointing every 10 seconds to ensure exactly-once processing guarantees and fault tolerance.
Running the Pipeline
Deploy the complete stack using the Make commands defined in the repository's README.
First, build and start the containers:
make up
This command builds a custom PyFlink image containing the Kafka and JDBC connectors, then launches the JobManager, TaskManager, and PostgreSQL services.
Next, submit the streaming job:
make job
This executes the PyFlink script inside the JobManager container, registering the tables and UDF before starting the continuous query.
Generate test events by visiting the demo endpoint, then query the results:
make psql
Inside the PostgreSQL shell, verify the enriched data:
SELECT ip, event_timestamp, geodata FROM processed_events LIMIT 5;
Summary
- Kafka Source: Configure the Kafka connector with SASL_SSL authentication to consume JSON events from Confluent Cloud, parsing timestamps with SQL functions.
- Python UDF: Implement the
GetLocationscalar function to enrich IP addresses with geolocation data via external API calls. - PostgreSQL Sink: Use the JDBC connector to sink transformed records into a relational database with exactly-once semantics.
- Containerized Deployment: Leverage Docker Compose to orchestrate Flink JobManager, TaskManager, and PostgreSQL services with environment-based configuration.
- Continuous Processing: Execute a persistent SQL query that streams data from Kafka through enrichment logic to the sink without batch windows.
Frequently Asked Questions
What is the difference between the Flink JobManager and TaskManager?
The JobManager coordinates the distributed execution of the streaming application, handling scheduling, checkpointing, and failure recovery. The TaskManager executes the specific subtasks defined in the job graph, such as reading from Kafka partitions or executing the Python UDF. In the DataExpert-io/data-engineer-handbook setup, the JobManager runs the start_job.py submission while TaskManager instances handle the parallel stream processing.
How does the pipeline handle authentication with Confluent Kafka?
The Kafka source table configures SASL_SSL security by embedding the sasl.jaas.config property directly in the SQL DDL. This property references environment variables for the Kafka API key and secret, allowing the Flink Kafka consumer to authenticate securely against Confluent Cloud without hardcoded credentials.
Can I modify the UDF to use a different geolocation service?
Yes. The GetLocation class in [start_job.py](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py#L58-L80) implements the eval method that makes HTTP requests to ip2location.io. You can replace the API endpoint and response parsing logic to integrate alternative services like MaxMind or IPStack, provided the function returns a JSON string matching the schema expected by the sink table.
How do I scale the pipeline to handle higher throughput?
Increase the parallelism by adjusting the taskmanager.numberOfTaskSlots configuration in the Docker Compose environment variables and ensuring the Kafka topic has sufficient partitions. The Flink cluster will automatically distribute the workload across available TaskManager slots, allowing concurrent processing of multiple Kafka partitions.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →