Building Real-Time Streaming Pipelines with Apache Flink and Kafka

Real-time streaming pipelines with Apache Flink and Kafka combine Kafka's high-throughput message ingestion with Flink's stateful stream processing to continuously transform and sink data to downstream systems.

This guide walks through a production-ready implementation from the DataExpert.io data-engineer-handbook repository that ingests web traffic events from Kafka, enriches them with geolocation data, and persists results to PostgreSQL—all running in Docker containers.

Architecture Overview

The pipeline follows a classic three-stage streaming architecture:

  1. Ingest: Kafka source connector pulls JSON events from a Confluent Cloud topic
  2. Process: Flink SQL with Python UDF enriches IP addresses with location metadata
  3. Sink: JDBC connector writes enriched records to PostgreSQL

According to the data-engineer-handbook source code, this pattern is implemented using PyFlink's Table API with SQL DDL statements for connector configuration.

Project Structure and Key Files

File Purpose Direct Link
README.md Setup instructions and Make targets View source
src/job/start_job.py Core PyFlink job with source, sink, and UDF definitions View source
docker-compose.yml Flink JobManager, TaskManager, and PostgreSQL services View source
example.env Environment variable template for credentials View source

The Docker Compose configuration provisions a complete Flink runtime. The JobManager exposes the Web UI on port 8081, while the TaskManager executes the streaming tasks.

Key configuration from [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):

  • JobManager: REST port 8081, bundled with Kafka connector JARs
  • TaskManager: 4 task slots for parallel execution
  • PostgreSQL: Port 5432 with persistent volume for the sink table
  • Environment injection: All Kafka and database credentials loaded from flink-env.env

Deploy the stack:


# Clone and navigate to the Flink training directory

git clone https://github.com/DataExpert-io/data-engineer-handbook.git
cd intermediate-bootcamp/materials/4-apache-flink-training

# Copy and fill environment template

cp example.env flink-env.env

# Edit flink-env.env with your Confluent Cloud credentials

# Build and start all services

make up

This executes docker compose --env-file flink-env.env up --build -d, which builds a custom Flink image with Python dependencies and the Kafka/JDBC connectors pre-installed.

Configuring the Kafka Source Table

The streaming pipeline begins with a Kafka source table defined in [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#L83-L100). This DDL creates a dynamic table that continuously polls the Kafka topic:

def create_events_source_kafka(t_env: StreamTableEnvironment, schema: str):
    """
    Creates a source table connected to our kafka topic, notice the parameters
    """
    table_name = "events"
    pattern = "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'"
    t_env.execute_sql(f"""
        CREATE TABLE {table_name} (
            url VARCHAR,
            referrer VARCHAR,
            user_agent VARCHAR,
            host VARCHAR,
            ip VARCHAR,
            headers VARCHAR,
            event_time VARCHAR,
            event_timestamp AS TO_TIMESTAMP(event_time, '{pattern}')
        ) 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'
        );
    """)
    return table_name

Critical configuration elements for Kafka source tables:

  • scan.startup.mode: latest-offset begins consuming new messages only; use earliest-offset for historical replay
  • event_timestamp: Computed column converts ISO-8601 strings to Flink TIMESTAMP for event-time processing
  • SASL/SSL: Required for Confluent Cloud and other secured Kafka deployments

Building the PostgreSQL Sink Table

The processed records land in a JDBC-connected PostgreSQL table. The sink definition 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):

def create_processed_events_sink_postgres(t_env: StreamTableEnvironment, schema: str):
    table_name = "processed_events"
    sink_ddl = f"""
        CREATE TABLE {table_name} (
            ip VARCHAR,
            event_timestamp TIMESTAMP(3),
            referrer VARCHAR,
            host VARCHAR,
            url VARCHAR,
            geodata VARCHAR,
            PRIMARY KEY (ip, event_timestamp) NOT ENFORCED
        ) WITH (
            'connector' = 'jdbc',
            'url' = '${{POSTGRES_URL}}',
            'table-name' = '{table_name}',
            'username' = '${{POSTGRES_USER}}',
            'password' = '${{POSTGRES_PW}}',
            'driver' = 'org.postgresql.Driver'
        );
    """
    t_env.execute_sql(sink_ddl)
    return table_name

Key sink parameters:

  • Primary key: (ip, event_timestamp) enables idempotent upserts—duplicate events with the same key overwrite previous values
  • NOT ENFORCED: Flink trusts the source data uniqueness; no runtime constraint checking overhead
  • PostgreSQL driver: Bundled in the custom Docker image

Enriching Streams with Python UDFs

The pipeline's enrichment logic calls an external IP geolocation service. The GetLocation UDF 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):

class GetLocation(ScalarFunction):
    """
    For every IP address, fetch the geolocation.
    """
    def eval(self, ip_address):
        response = requests.get("https://api.ip2location.io/", params={
            "ip": ip_address,
            "key": os.environ.get("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", ""),
        })

# Register the UDF for use in SQL

get_location = udf(GetLocation(), result_type=DataTypes.STRING())

UDF performance considerations:

  • ScalarFunction: Processes one row at a time; for high-throughput scenarios, consider TableFunction for batch lookups or async external service calls
  • Error handling: Returns empty JSON on API failure to prevent job failures from transient service outages
  • Environment variables: API keys injected at runtime, never hardcoded

Orchestrating the Streaming Job

The log_processing function 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#L103-L126) wires components together and executes the continuous INSERT:

def log_processing():
    # Load environment variables

    env = StreamExecutionEnvironment.get_execution_environment()
    env.enable_checkpointing(10 * 1000)  # 10-second checkpoints for fault tolerance

    settings = EnvironmentSettings.new_instance().in_streaming_mode().build()
    t_env = StreamTableEnvironment.create(env, settings)

    # Register the UDF

    t_env.create_temporary_system_function("get_location", get_location)

    # Create source and sink tables

    source_table = create_events_source_kafka(t_env, "test")
    postgres_sink = create_processed_events_sink_postgres(t_env, "test")

    # Execute continuous streaming query

    t_env.execute_sql(f"""
        INSERT INTO {postgres_sink}
        SELECT ip, event_timestamp, referrer, host, url, get_location(ip) AS geodata
        FROM {source_table}
    """)

Checkpointing every 10 seconds ensures exactly-once processing semantics without excessive latency.

Running and Monitoring the Pipeline

Submit the job using the provided Make target:


# Submit PyFlink job to running cluster

make job

# Equivalent manual command:

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

Monitor execution:

Command Purpose
make ui Open Flink Web Dashboard at http://localhost:8081
make psql Interactive PostgreSQL shell to query results
make down Stop all containers and clean up

To generate test events, visit https://bootcamp.techcreator.io/—each page interaction publishes a JSON message to the Kafka topic that Flink immediately processes.

Querying Enriched Results

After the job runs, verify data in PostgreSQL:

make psql
-- View enriched web traffic with geolocation
SELECT 
    ip,
    event_timestamp,
    url,
    geodata::json->>'country' as country,
    geodata::json->>'city' as city
FROM processed_events
ORDER BY event_timestamp DESC
LIMIT 10;

Sample output:


        ip        |     event_timestamp      |           url            | country |    city
------------------+--------------------------+--------------------------+---------+------------
 203.0.113.42     | 2024-01-15 09:23:17.234  | /courses/flink-basics    | US      | San Francisco
 198.51.100.15    | 2024-01-15 09:22:58.891  | /blog/streaming-patterns | DE      | Berlin

Extending the Pipeline

The data-engineer-handbook implementation demonstrates patterns applicable to production workloads:

  • Multiple sources: Add additional Kafka topics or CDC connectors (Debezium for PostgreSQL/MySQL)
  • Windowed aggregations: Replace INSERT … SELECT with TUMBLE or HOP windows for real-time analytics
  • Exactly-once sinks: The JDBC connector with primary keys provides idempotent writes; for Kafka-to-Kafka pipelines, use Flink's Kafka producer with transactional IDs
  • Schema evolution: Flink's JSON format supports field addition; use ROW<> types for nested structures

Summary

  • Apache Flink and Kafka integrate through SQL DDL connectors that treat streaming topics as dynamic tables
  • PyFlink enables Python UDFs for external service calls while maintaining Flink's checkpointing and exactly-once guarantees
  • Docker Compose local development mirrors production deployment patterns with configurable task parallelism
  • Idempotent sinks with primary keys prevent duplicate data when jobs restart from checkpoints

Frequently Asked Questions

Ten seconds balances recovery speed with processing overhead for most pipelines. Lower intervals (1-5 seconds) reduce data replay after failures but increase storage and CPU costs. For high-throughput, low-latency applications, consider incremental checkpoints and local state backend tuning.

The JSON format deserializer ignores unknown fields by default, enabling backward-compatible schema evolution. For structured schema management, integrate with Confluent Schema Registry using the avro-confluent format and specify 'properties.schema.registry.url'.

Yes—Flink supports rescaling through savepoints. Trigger a savepoint, cancel the job, adjust parallelism in your submission configuration, and restart from the savepoint. The Kafka connector resumes from committed offsets transparently.

PyFlink provides accessibility for Python-native data teams and seamless integration with Python ML libraries (pandas, scikit-learn). The trade-off is slight serialization overhead for UDF execution. For purely SQL pipelines with no custom Python logic, Flink SQL (via SQL Client or Table API) offers identical performance across all language bindings.

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 →