Building Real‑Time Streaming Pipelines with Apache Flink and Kafka: A Complete Architecture Guide

A real‑time streaming pipeline with Apache Flink and Kafka ingests events from Kafka topics, enriches them through user‑defined functions, and writes results to external sinks like PostgreSQL using declarative SQL DDL and continuous INSERT statements.

The DataExpert-io/data-engineer-handbook repository demonstrates production‑grade patterns for low‑latency stream processing. This guide breaks down the architecture implemented in the Apache Flink training module, showing how to wire together Kafka sources, Python UDFs, and JDBC sinks into a fault‑tolerant, containerized pipeline.

Architecture Overview: The Three‑Layer Pipeline

The pipeline follows a standard extract‑transform‑load pattern optimized for streaming. Raw web traffic events flow from a Confluent Kafka cluster into Flink, where a Python UDF augments each record with geolocation metadata, and the enriched stream lands in PostgreSQL.

According to the source code, the log_processing function orchestrates three distinct phases:

  1. Source: A Kafka connector table that deserializes JSON and converts event time strings to timestamps.
  2. Transform: A scalar UDF GetLocation that queries the ip2location.io API.
  3. Sink: A PostgreSQL connector table using the JDBC driver for reliable upserts.

Configuring the Kafka Source Table

Flink’s Kafka connector supports SASL_SSL authentication and consumer groups for production workloads. The DDL 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#L83-L100) defines the source with explicit timestamp parsing:

CREATE TABLE events (
    url VARCHAR,
    referrer VARCHAR,
    user_agent VARCHAR,
    host VARCHAR,
    ip VARCHAR,
    headers VARCHAR,
    event_time VARCHAR,
    event_timestamp AS TO_TIMESTAMP(event_time, 
        'yyyy-MM-dd''T''HH:mm:ss.SSS''Z''')
) WITH (
    'connector' = 'kafka',
    'properties.bootstrap.servers' = '${KAFKA_URL}',
    'topic' = '${KAFKA_TOPIC}',
    'properties.group.id' = '${KAFKA_GROUP}',
    'properties.security.protocol' = 'SASL_SSL',
    'properties.sasl.mechanism' = 'PLAIN',
    'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="${KAFKA_WEB_TRAFFIC_KEY}" password="${KAFKA_WEB_TRAFFIC_SECRET}";',
    'scan.startup.mode' = 'latest-offset',
    'format' = 'json'
)

This configuration leverages Flink’s SQL variable substitution to inject credentials securely via environment variables defined in example.env.

Enriching Data with Python UDFs

Real‑time enrichment moves computation to the stream rather than the database. The GetLocation class implemented 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) extends Flink’s ScalarFunction to perform synchronous HTTP lookups:

class GetLocation(ScalarFunction):
    def eval(self, ip_address):
        response = requests.get(
            "https://api.ip2location.io",
            params={'ip': ip_address, 'key': os.getenv("IP_CODING_KEY")}
        )
        if response.status_code != 200:
            return json.dumps({})
        data = json.loads(response.text)
        return json.dumps({
            "country": data.get('country_code', ''),
            "state": data.get('region_name', ''),
            "city": data.get('city_name', '')
        })

get_location = udf(GetLocation(), result_type=DataTypes.STRING())
t_env.register_function("get_location", get_location)

The UDF returns a JSON string containing country, state, and city codes, which the SQL query projects as a computed column during the INSERT operation.

Writing Results to PostgreSQL

The sink table uses Flink’s JDBC connector with PostgreSQL‑specific dialect support. Defined 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#L36-L56), the DDL configures the connection URL and driver:

CREATE TABLE processed_events (
    ip VARCHAR,
    event_timestamp TIMESTAMP(3),
    referrer VARCHAR,
    host VARCHAR,
    url VARCHAR,
    geodata VARCHAR
) WITH (
    'connector' = 'jdbc',
    'url' = '${POSTGRES_URL}',
    'driver' = 'org.postgresql.Driver',
    'table-name' = 'processed_events',
    'username' = '${POSTGRES_USER}',
    'password' = '${POSTGRES_PASSWORD}'
)

Flink handles connection pooling and retry logic automatically, ensuring exactly‑once semantics when checkpointing is enabled.

Orchestrating Infrastructure with Docker Compose

The [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) defines a three‑service architecture:

  • JobManager: Exposes port 8081 for the Web UI and accepts PyFlink job submissions.
  • TaskManager: Runs with configurable slots (taskmanager.numberOfTaskSlots: 10) to parallelize the UDF execution.
  • PostgreSQL: Hosts the destination table with persistent volume storage.

Environment variables flow through flink-env.env, keeping secrets out of the image layers.

Running the Streaming Job

The repository uses Make targets to abstract Docker commands. Execute the pipeline locally with:


# Build images and start services

make up

# Submit the PyFlink job to the JobManager

make job

# Generate test events (visit the demo site)

open https://bootcamp.techcreator.io/

# Query enriched results

make psql

After submission, navigate to http://localhost:8081/ to view the running job graph. The continuous SQL query executes indefinitely, checkpointing every ten seconds to guarantee fault tolerance:

table_env.execute_sql("""
    INSERT INTO processed_events
    SELECT
        ip,
        event_timestamp,
        referrer,
        host,
        url,
        get_location(ip) AS geodata
    FROM events
""")

Summary

  • Declarative Sources/Sinks: Flink SQL DDL handles connector configuration, schema mapping, and authentication for both Kafka and PostgreSQL without imperative boilerplate.
  • UDF Enrichment: Python scalar functions integrate external APIs directly into the streaming query, minimizing latency between ingestion and augmentation.
  • Containerized Deployment: Docker Compose bundles the JobManager, TaskManager, and database with environment‑variable injection for reproducible local development.
  • Fault Tolerance: Checkpointing every ten seconds ensures exactly‑once processing semantics even during container restarts or task failures.

Frequently Asked Questions

Flink supports SASL_SSL through the properties.sasl.jaas.config option in the Kafka connector DDL. The configuration string embeds username and password variables that resolve at runtime from environment variables injected via the Docker Compose env_file, as shown in the events table definition.

What performance considerations apply to the GetLocation UDF?

Because GetLocation performs synchronous HTTP requests inside the eval method, it introduces network latency per record. The TaskManager configuration in [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) allocates ten task slots to parallelize these calls across multiple threads, preventing head‑of‑line blocking for high‑volume streams.

How do you scale this pipeline for production workloads?

Increase the taskmanager.numberOfTaskSlots value in the Docker Compose configuration to add parallelism, or deploy additional TaskManager containers behind a Flink session cluster. For Kafka source scaling, increase the KAFKA_GROUP consumer count and ensure the topic has sufficient partitions to match Flink’s parallelism.

How can you verify the pipeline is processing events in real time?

After running make job, visit the demo site to generate Kafka messages, then execute make psql to query the processed_events table. If the architecture is functioning correctly, the PostgreSQL table will contain rows with populated geodata JSON within seconds of the event timestamp, confirming end‑to‑end latency through Flink’s streaming engine.

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 →