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

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.

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, 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
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

The core processing logic resides in 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.

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.

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.

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 to create the target schema.

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 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
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 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 for the pipeline logic and sql/init.sql for schema definition.

Frequently Asked Questions

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 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, 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. 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.

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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →