Jaehyeon Kim
Jaehyeon Kim

  • Blog
    • Categories

      List of categories.

    • Tags

      List of tags.

    • Series

      List of series.

    • Archives

  • Projects
  • Slides

/

  • Github Linkedin RSS

  • Font Size
  • Palette
  • Mode
  1. Home
  2. Tags
  3. Apache Beam

Develop Streaming File Reader using Splittable DoFn - Apache Beam Python Examples Part 10

Develop Streaming File Reader using Splittable DoFn - Apache Beam Python Examples Part 10
December 19, 202412 min read Data-Streaming Apache Beam Python ExamplesApache BeamApache FlinkPythonSplittable DoFn

A streaming file reader built with Splittable DoFn scans an input folder for new files repeatedly, a pattern for unbounded sources in the Python SDK.

Read More: Develop Streaming File Reader using Splittable DoFn - Apache Beam Python Examples Part 10

Develop Batch File Reader and PiSampler using Splittable DoFn - Apache Beam Python Examples Part 9

Develop Batch File Reader and PiSampler using Splittable DoFn - Apache Beam Python Examples Part 9
December 5, 202410 min read Data-Streaming Apache Beam Python ExamplesApache BeamApache FlinkPythonSplittable DoFn

Splittable DoFn in Beam Python builds a batch file reader that processes files in parallel, and a PiSampler that estimates pi by Monte Carlo runs.

Read More: Develop Batch File Reader and PiSampler using Splittable DoFn - Apache Beam Python Examples Part 9

Enhance Sport Activity Tracker with Runner Motivation - Apache Beam Python Examples Part 8

Enhance Sport Activity Tracker with Runner Motivation - Apache Beam Python Examples Part 8
November 21, 202418 min read Data-Streaming Apache Beam Python ExamplesApache BeamApache FlinkApache KafkaPython

Pacing messages are added to the Beam sport activity tracker by comparing short term speed metrics against their long term counterparts.

Read More: Enhance Sport Activity Tracker with Runner Motivation - Apache Beam Python Examples Part 8

Apache Beam Python Examples - Part 7 Separate Droppable Data into Side Output

Apache Beam Python Examples - Part 7 Separate Droppable Data into Side Output
October 24, 202418 min read Data-Streaming Apache Beam Python ExamplesApache BeamApache FlinkApache KafkaPython

Late droppable elements are detected by a timer in a stateful DoFn and sent to a Beam side output instead of being discarded without notice.

Read More: Apache Beam Python Examples - Part 7 Separate Droppable Data into Side Output

Call RPC Service in Batch with Defined Batch Size using Stateful DoFn - Apache Beam Python Examples Part 6

Call RPC Service in Batch with Defined Batch Size using Stateful DoFn - Apache Beam Python Examples Part 6
October 2, 202414 min read Data-Streaming Apache Beam Python ExamplesApache BeamApache FlinkApache KafkaGRPCPython

A stateful DoFn with Beam state and timers fixes the gRPC batch size and maximum wait time instead of leaving the bundle size to the runner.

Read More: Call RPC Service in Batch with Defined Batch Size using Stateful DoFn - Apache Beam Python Examples Part 6

Call RPC Service in Batch using Stateless DoFn - Apache Beam Python Examples Part 5

Call RPC Service in Batch using Stateless DoFn - Apache Beam Python Examples Part 5
September 18, 202412 min read Data-Streaming Apache Beam Python ExamplesApache BeamApache FlinkApache KafkaGRPCPython

Batching gRPC calls in a stateless DoFn so one request covers a whole bundle, cutting the time a Beam Python pipeline spends on enrichment.

Read More: Call RPC Service in Batch using Stateless DoFn - Apache Beam Python Examples Part 5

Cache Data on Apache Beam Pipelines Using a Shared Object

Cache Data on Apache Beam Pipelines Using a Shared Object
August 22, 20248 min read Data ProcessingApache BeamCachingData EnrichmentPython

The Shared class in the Beam Python SDK caches lookup data in memory for batch and streaming pipelines, with a periodic refresh for the latter.

Read More: Cache Data on Apache Beam Pipelines Using a Shared Object

Apache Beam Python Examples - Part 4 Call RPC Service for Data Augmentation

Apache Beam Python Examples - Part 4 Call RPC Service for Data Augmentation
August 15, 202414 min read Data-Streaming Apache Beam Python ExamplesApache BeamApache FlinkApache KafkaGRPCPython

Data augmentation in Beam Python by calling a gRPC service once per input element, running on a local Flink cluster with Kafka as the source.

Read More: Apache Beam Python Examples - Part 4 Call RPC Service for Data Augmentation

Build Sport Activity Tracker with/without SQL - Apache Beam Python Examples Part 3

Build Sport Activity Tracker with/without SQL - Apache Beam Python Examples Part 3
August 1, 202420 min read Data-Streaming Apache Beam Python ExamplesApache BeamApache FlinkApache KafkaPython

A sport activity tracker in Beam Python, built first with native transforms and then with Beam SQL, showing the limits of Beam SQL in the Python SDK.

Read More: Build Sport Activity Tracker with/without SQL - Apache Beam Python Examples Part 3

Calculate Average Word Length with/without Fixed Look back - Apache Beam Python Examples Part 2

Calculate Average Word Length with/without Fixed Look back - Apache Beam Python Examples Part 2
July 18, 202415 min read Data-Streaming Apache Beam Python ExamplesApache BeamApache FlinkApache KafkaPython

Two Beam Python pipelines compute average word length from a Kafka topic, one emitting a global average and one using a sliding time window.

Read More: Calculate Average Word Length with/without Fixed Look back - Apache Beam Python Examples Part 2
  • ««
  • «
  • 1
  • 2
  • »
  • »»
Profile
Jaehyeon Kim
Jaehyeon Kim
Data Engineer | Data Streaming | Powering ML & AI in Real Time
Taxonomies
Data Streaming 70 Data Engineering 42 Development 28 Data Analysis 17 Data Integration 12 Open Source 12 Machine Learning 7 Kubernetes 6 Security 5 Data Architecture 3 Data Processing 3 Big Data 2 System Architecture 2 Web Development 2
Python 76 Apache Kafka 72 AWS 50 Docker 50 Apache Flink 38 Apache Beam 17 Apache Spark 16 Kafka Connect 16 Amazon MSK 14 AWS Lambda 14 Benchtop 13 dbt 13 Amazon EMR 11 dynamic-des 11 odctl 11 PostgreSQL 10 Kubernetes 8 PyFlink 8 Change Data Capture (CDC) 7 Debezium 7 Kotlin 7 Apache Airflow 6 Apache Iceberg 6 Discrete Event Simulation 6 Kafka UI 6 Amazon DynamoDB 5 PySpark 5 SimPy 5 Amazon API Gateway 4 Amazon Athena 4 AWS Glue 4 AWS Glue Schema Registry 4 BigQuery 4 Digital Twin 4 FastAPI 4 Minikube 4 Amazon EKS 3 Amazon QuickSight 3 Amazon S3 3 Apache Hudi 3 ALL 143
Kafka Development with Docker 11 Apache Beam Python Examples 10 Real Time Streaming with Kafka and Flink 7 Building Real-Time Digital Twins with dynamic-des 6 dbt Pizza Shop Demo 6 Apache Beam Local Development with Python 5 dbt for Effective Data Transformation on AWS 5 Getting Started with Real-Time Streaming in Kotlin 5 Kafka Connect for AWS Services Integration 5 Serverless Data Product 4 Data Lake Demo Using Change Data Capture 3 Getting Started with PyFlink on AWS 3 Kafka Development on Kubernetes 3 Parallel processing on single machine 3 Realtime Dashboard with FastAPI, Streamlit and Next.js 3 API development with R 2 dbt Guide for Production 2 Deploy Python Stream Processing App on Kubernetes 2 Download Stock Data 2 From Prototype to Production: Real-Time Product Recommendation with Contextual Bandits 2 ALL 25
2026 16 2025 13 2024 29 2023 39 2022 15 2021 7 2020 1 2019 5 2018 2 2017 6 2016 6 2015 15 2014 5
Posts
  • Defining Data-Streaming Simulations in YAML, Without Writing Python
    Defining Data-Streaming Simulations in YAML, Without Writing Python
    October 6, 2026
  • Keeping Game Leaderboards Up to Date in Real Time with Kafka and Flink SQL
    Keeping Game Leaderboards Up to Date in Real Time with Kafka and Flink SQL
    October 2, 2026
  • Change Data Capture on a Simulated Online Shop with Debezium and Kafka Connect
    Change Data Capture on a Simulated Online Shop with Debezium and Kafka Connect
    October 1, 2026
  • Building an Agentic Analytics System over an Iceberg Lakehouse
    Building an Agentic Analytics System over an Iceberg Lakehouse
    July 18, 2026
  • Current London 2026: Building End-to-End Data Lineage
    Current London 2026: Building End-to-End Data Lineage
    May 22, 2026
  • Building a Real-Time Industrial Digital Twin with Apache Flink and Online Machine Learning
    Building a Real-Time Industrial Digital Twin with Apache Flink and Online Machine Learning
    April 21, 2026
  • Productionizing an Online Product Recommender using Event Driven Architecture
    Productionizing an Online Product Recommender using Event Driven Architecture
    February 23, 2026
  • Stream Processing with Flink in Kotlin
    Stream Processing with Flink in Kotlin
    December 10, 2025
  • Guide to Building Integrated Web Applications with FastAPI and NiceGUI
    Guide to Building Integrated Web Applications with FastAPI and NiceGUI
    November 19, 2025
  • Self-service Data Platform via a Multi-tenant SQL Gateway
    Self-service Data Platform via a Multi-tenant SQL Gateway
    July 17, 2025
  • Defining Data-Streaming Simulations in YAML, Without Writing Python
    Defining Data-Streaming Simulations in YAML, Without Writing Python
    October 6, 2026
  • Keeping Game Leaderboards Up to Date in Real Time with Kafka and Flink SQL
    Keeping Game Leaderboards Up to Date in Real Time with Kafka and Flink SQL
    October 2, 2026
  • Change Data Capture on a Simulated Online Shop with Debezium and Kafka Connect
    Change Data Capture on a Simulated Online Shop with Debezium and Kafka Connect
    October 1, 2026
  • Data Streaming and Machine Learning Projects That Run on Your Laptop
    Data Streaming and Machine Learning Projects That Run on Your Laptop
    September 30, 2026
  • Learning MLOps with a Feature Store: A New Series
    Learning MLOps with a Feature Store: A New Series
    September 28, 2026
  • Building an Agentic Analytics System over an Iceberg Lakehouse
    Building an Agentic Analytics System over an Iceberg Lakehouse
    July 18, 2026
  • Dynamic DES: A Declarative API with Postgres and Redis Connectors
    Dynamic DES: A Declarative API with Postgres and Redis Connectors
    July 17, 2026
  • Running Kafka, Flink, Spark, Trino and Iceberg Locally with One CLI
    Running Kafka, Flink, Spark, Trino and Iceberg Locally with One CLI
    July 16, 2026
  • One Simulation, Two Pipelines: Batch Training and Live Inference with Dynamic DES
    One Simulation, Two Pipelines: Batch Training and Live Inference with Dynamic DES
    May 25, 2026
  • Current London 2026: Building End-to-End Data Lineage
    Current London 2026: Building End-to-End Data Lineage
    May 22, 2026
Actions
Go back Reload Copy URL

Jaehyeon Kim

Data Engineer | Data Streaming | Powering ML & AI in Real Time

Copyright © 2023-2026 Jaehyeon Kim. All Rights Reserved.