# How to Implement Real-Time Streaming with Kafka and Flink: A Complete PyFlink Tutorial

> Master real-time streaming with Kafka and Flink using PyFlink SQL. Learn to enrich data with UDFs and sink to PostgreSQL in this complete tutorial.

- Repository: [DataExpert.io/data-engineer-handbook](https://github.com/DataExpert-io/data-engineer-handbook)
- Tags: tutorial
- Published: 2026-08-08

---

**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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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/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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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/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/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/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](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/README.md).

First, build and start the containers:

```bash
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:

```bash
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:

```bash
make psql

```

Inside the PostgreSQL shell, verify the enriched data:

```sql
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 `GetLocation` scalar 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`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/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/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.