# How to Build Real-Time Streaming Pipelines with Apache Flink and Kafka: A PyFlink Tutorial

> Build real-time streaming pipelines using PyFlink and Kafka. Learn to configure the Table API, apply session windowing, and store results in PostgreSQL with this Docker Compose tutorial.

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

---

**You can build real-time streaming pipelines by configuring PyFlink's Table API to consume from Kafka with SASL authentication, apply session windowing for aggregation, and persist results to PostgreSQL via JDBC, all orchestrated through Docker Compose as implemented in the DataExpert-io/data-engineer-handbook.**

The **DataExpert-io/data-engineer-handbook** repository provides a production-ready reference architecture demonstrating how to build real-time streaming pipelines with Apache Flink and Kafka. This intermediate bootcamp material leverages PyFlink, Docker Compose, and Confluent Cloud to create an end-to-end data pipeline that sessionizes web traffic events and stores them in PostgreSQL for downstream analytics.

## Architecture of the Flink and Kafka Pipeline

The pipeline architecture follows a standard Lambda-style ingestion pattern with three primary layers. **Kafka** acts as the distributed message broker, streaming web-traffic events from a Confluent Cloud cluster. **Apache Flink** processes these events using PyFlink's Table API to perform sessionization—grouping events by IP address and host within 5-minute session windows. Finally, **PostgreSQL** serves as the persistent sink, storing aggregated results in the `processed_events` table.

According to the source code in `intermediate-bootcamp/materials/4-apache-flink-training/`, all components run as containerized services defined in [`docker-compose.yml`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/docker-compose.yml), including the Flink JobManager, TaskManager, and PostgreSQL instance.

## Environment Configuration

Before running the pipeline, you must configure connection credentials and runtime parameters. The repository provides `example.env` as a template that you copy to `flink-env.env`.

Key configuration variables include:

- **KAFKA_URL**: The Confluent Cloud bootstrap server (e.g., `pkc-rgm37.us-west-2.aws.confluent.cloud:9092`)
- **KAFKA_TOPIC**: The source topic name (`bootcamp-events-prod`)
- **KAFKA_WEB_TRAFFIC_KEY** and **KAFKA_WEB_TRAFFIC_SECRET**: SASL authentication credentials
- **POSTGRES_URL**: JDBC connection string pointing to the containerized database (`jdbc:postgresql://host.docker.internal:5432/postgres`)
- **FLINK_VERSION**: Set to `1.16.0` with **PYTHON_VERSION** `3.7.9`

```text
KAFKA_WEB_TRAFFIC_SECRET="<GET FROM bootcamp.techcreator.io>"
KAFKA_WEB_TRAFFIC_KEY="<GET FROM bootcamp.techcreator.io>"
KAFKA_GROUP=web-events
KAFKA_TOPIC=bootcamp-events-prod
KAFKA_URL=pkc-rgm37.us-west-2.aws.confluent.cloud:9092

FLINK_VERSION=1.16.0
PYTHON_VERSION=3.7.9

POSTGRES_URL="jdbc:postgresql://host.docker.internal:5432/postgres"
POSTGRES_USER=postgres
POSTGRES_PASSWORD=postgres
POSTGRES_DB=postgres

```

## Implementing the PyFlink Job

The core processing logic resides in [`intermediate-bootcamp/materials/4-apache-flink-training/src/job/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). This script uses the **Table API** to define the Kafka source, apply watermarks for event-time processing, and execute a session window aggregation before writing to the JDBC sink.

### Defining the Kafka Source

The job creates a temporary table named `KafkaSource` using the Kafka connector. It configures **SASL_SSL** security with the PLAIN mechanism, passing the Confluent Cloud credentials via the `sasl.jaas.config` option.

```python
from pyflink.table import StreamTableEnvironment, EnvironmentSettings
from pyflink.table import TableDescriptor, Schema, DataTypes

env_settings = EnvironmentSettings.in_streaming_mode()
t_env = StreamTableEnvironment.create(environment_settings=env_settings)

kafka_source = TableDescriptor.for_connector("kafka") \
    .schema(Schema.new_builder()
        .column("ip", DataTypes.STRING())
        .column("host", DataTypes.STRING())
        .column("timestamp", DataTypes.TIMESTAMP(3))
        .column("event", DataTypes.STRING())
        .watermark("timestamp", "timestamp - INTERVAL '5' SECOND")
        .build()) \
    .option("topic", "${KAFKA_TOPIC}") \
    .option("properties.bootstrap.servers", "${KAFKA_URL}") \
    .option("properties.security.protocol", "SASL_SSL") \
    .option("properties.sasl.mechanism", "PLAIN") \
    .option("properties.sasl.jaas.config",
            "org.apache.kafka.common.security.plain.PlainLoginModule required "
            "username='${KAFKA_WEB_TRAFFIC_KEY}' password='${KAFKA_WEB_TRAFFIC_SECRET}';") \
    .format("json") \
    .build()

t_env.create_temporary_table("KafkaSource", kafka_source)

```

### Session Window Aggregation

The pipeline implements **sessionization** by defining a session window with a 5-minute gap. Events from the same IP and host that occur within 5 minutes of each other are grouped into a single session.

```python
from pyflink.table.expressions import col, session

events = t_env.from_path("KafkaSource")

sessionized = events.window(
    session(col("timestamp"), lit(5).minutes)
).group_by(
    col("ip"), col("host")
).select(
    col("ip"),
    col("host"),
    col("event").count.alias("cnt"),
    col("timestamp").max.alias("latest_ts")
)

```

### Configuring the PostgreSQL Sink

The processed results are written to a JDBC sink using the PostgreSQL driver. The sink table schema must match the aggregation output.

```python
postgres_sink = TableDescriptor.for_connector("jdbc") \
    .schema(Schema.new_builder()
        .column("ip", DataTypes.STRING())
        .column("host", DataTypes.STRING())
        .column("cnt", DataTypes.BIGINT())
        .column("latest_ts", DataTypes.TIMESTAMP(3))
        .build()) \
    .option("url", "${POSTGRES_URL}") \
    .option("table-name", "processed_events") \
    .option("driver", "org.postgresql.Driver") \
    .option("username", "${POSTGRES_USER}") \
    .option("password", "${POSTGRES_PASSWORD}") \
    .build()

t_env.create_temporary_table("PostgresSink", postgres_sink)
sessionized.execute_insert("PostgresSink")

```

## Database Schema Setup

Before submitting the Flink job, you must initialize the PostgreSQL sink table. The repository includes [`intermediate-bootcamp/materials/4-apache-flink-training/sql/init.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/sql/init.sql) to create the target schema.

```sql
CREATE TABLE IF NOT EXISTS processed_events (
    ip VARCHAR,
    host VARCHAR,
    cnt BIGINT,
    latest_ts TIMESTAMP
);

```

Run this script against the containerized PostgreSQL instance to ensure the sink table exists before the job starts writing data.

## Docker Compose Orchestration

The [`docker-compose.yml`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/docker-compose.yml) file defines four critical services: `jobmanager`, `taskmanager`, `postgres`, and the Kafka broker configuration. The Flink services mount the `src/` directory to access the PyFlink script at runtime.

### Makefile Automation

The repository includes a **Makefile** to simplify pipeline operations:

- **`make up`**: Builds the Flink base image (including PyFlink and Kafka connectors) and starts the full stack
- **`make job`**: Submits the PyFlink job to the JobManager using `flink run -py /opt/src/job/start_job.py -d`
- **`make psql`**: Opens an interactive PostgreSQL shell to query the `processed_events` table
- **`make stop`**: Gracefully stops the Flink cluster and supporting services

```makefile
build:
	docker build -t flink-base .

up: build
	docker-compose up -d

job:
	docker-compose exec jobmanager ./bin/flink run -py /opt/src/job/start_job.py -d

psql:
	docker exec -it postgres psql -U postgres -d postgres

clean:
	docker-compose down -v

```

## Running the End-to-End Pipeline

Follow these steps to execute the real-time streaming pipeline:

1. **Initialize environment variables**: Copy `example.env` to `flink-env.env` and populate your Confluent Cloud credentials and IP geolocation API key.

2. **Start the infrastructure**: Run `make up` to build the custom Flink image and launch the JobManager, TaskManager, and PostgreSQL containers.

3. **Create the sink table**: Execute [`sql/init.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/sql/init.sql) against the PostgreSQL container to create the `processed_events` table.

4. **Submit the Flink job**: Execute `make job` to deploy the PyFlink application. The job will begin consuming from the `bootcamp-events-prod` Kafka topic.

5. **Generate test events**: Visit `https://bootcamp.techcreator.io/` to trigger web traffic events that publish to the Kafka topic.

6. **Verify results**: Run `make psql` and query the `processed_events` table to confirm that sessionized aggregates are persisting correctly.

7. **Teardown**: Use `make clean` to stop all containers and remove volumes when finished testing.

## Summary

- **Apache Flink** processes real-time streams using the Table API with event-time semantics and watermarking.
- **Kafka** integration requires SASL_SSL configuration for secure authentication against Confluent Cloud.
- **Session windows** group related events (by IP and host) using a 5-minute inactivity gap before emitting aggregates.
- **PostgreSQL** serves as the JDBC sink for persistent storage of processed session data.
- **Docker Compose** and the provided **Makefile** automate the deployment of JobManagers, TaskManagers, and database services.
- Source files are located in `intermediate-bootcamp/materials/4-apache-flink-training/`, specifically [`src/job/start_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/src/job/start_job.py) for the pipeline logic and [`sql/init.sql`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/sql/init.sql) for schema definition.

## Frequently Asked Questions

### What is the difference between the Flink JobManager and TaskManager in this architecture?

The **JobManager** coordinates the distributed execution of the PyFlink job, handling scheduling, checkpointing, and recovery, while **TaskManagers** execute the actual data processing tasks (reading from Kafka, windowing, writing to PostgreSQL). In the [`docker-compose.yml`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/docker-compose.yml) configuration, the JobManager receives the job submission via `make job`, then distributes work across available TaskManager slots.

### How does the pipeline handle Kafka authentication securely?

The pipeline uses **SASL_SSL** with the PLAIN mechanism, passing credentials through the `sasl.jaas.config` property in the Kafka connector configuration. As shown in [`src/job/start_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/src/job/start_job.py), these values reference environment variables injected via Docker Compose from the `flink-env.env` file, keeping sensitive keys out of the source code.

### Why use a 5-minute gap for session windows?

The 5-minute session gap groups web traffic events into meaningful user sessions—if a user (identified by IP and host) generates events within 5 minutes of each other, Flink treats them as a single session. This aggregation reduces data volume and provides behavioral analytics (page views per session) rather than processing raw, individual events.

### Can I modify this pipeline to write to a different sink instead of PostgreSQL?

Yes, you can replace the JDBC sink with any supported connector by modifying the `TableDescriptor` in [`start_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/start_job.py). Flink supports sinks to Elasticsearch, S3, or other Kafka topics by changing the connector type from `"jdbc"` to `"elasticsearch-7"` or `"filesystem"` and updating the corresponding options, though you must ensure the necessary JAR dependencies are included in the Flink container image.