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.
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, 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.0with PYTHON_VERSION3.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
Implementing the PyFlink Job
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 stackmake job: Submits the PyFlink job to the JobManager usingflink run -py /opt/src/job/start_job.py -dmake psql: Opens an interactive PostgreSQL shell to query theprocessed_eventstablemake 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:
-
Initialize environment variables: Copy
example.envtoflink-env.envand populate your Confluent Cloud credentials and IP geolocation API key. -
Start the infrastructure: Run
make upto build the custom Flink image and launch the JobManager, TaskManager, and PostgreSQL containers. -
Create the sink table: Execute
sql/init.sqlagainst the PostgreSQL container to create theprocessed_eventstable. -
Submit the Flink job: Execute
make jobto deploy the PyFlink application. The job will begin consuming from thebootcamp-events-prodKafka topic. -
Generate test events: Visit
https://bootcamp.techcreator.io/to trigger web traffic events that publish to the Kafka topic. -
Verify results: Run
make psqland query theprocessed_eventstable to confirm that sessionized aggregates are persisting correctly. -
Teardown: Use
make cleanto 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/, specificallysrc/job/start_job.pyfor the pipeline logic andsql/init.sqlfor 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 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →