尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

PyFlink 流处理实战:基于 Redpanda、Flink 与 Postgres 的出租车数据会话窗口作业

PyFlink 流处理实战:基于 Redpanda、Flink 与 Postgres 的出租车数据会话窗口作业 PyFlink 流处理实战基于 Redpanda、Flink 与 Postgres 的出租车数据会话窗口作业【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp导读本篇技术指南以 Data Engineering Zoomcamp 2027 届07-streaming模块中 PyFlink 作业练习 为核心完整还原一套基于RedpandaKafka 协议兼容 Apache Flink PostgreSQL的真实流处理实验先用 Python 生产者把 2019 年 10 月绿色出租车Green Taxi数据写入 Kafka topic再编写 PyFlink 作业消费数据、建立会话窗口Session Window统计连续运营时段。读完本文你将掌握本地 Docker 化流环境搭建、kafka-python 生产端、PyFlink 表 API 的 DDL 建表、事件时间与水位线Watermark配置以及会话窗口聚合的完整实战方案。一、作业背景与数据源本次作业使用纽约绿色出租车 2019 年 10 月的数据green_tripdata_2019-10.csv.gz数据来自 DataTalksClub 的 nyc-tlc-data 发布。该数据集中每条记录代表一次出租车行程本次练习只关心以下字段lpep_pickup_datetime上车时间lpep_dropoff_datetime下车时间PULocationID上车地点 IDDOLocationID下车地点 IDpassenger_count乘客数trip_distance行程距离tip_amount小费金额整体实验分为四个阶段对应作业中的若干问题启动 Redpanda、Flink JobManager、Flink TaskManager 与 Postgres 四类服务验证 Python 客户端能连接 Kafka 服务器把出租车数据写入green-tripstopic编写带会话窗口的 PyFlink 作业找出连续运营无 5 分钟以上中断最久的上下车地点组合。二、启动流处理环境Redpanda Flink Postgres作业目录 extras/pyflink 下已提供 docker-compose.yml将其复制到你的作业目录后直接启动即可docker-compose up # 需要后台运行时加 -ddocker-compose up -d2.1 四个核心服务从 docker-compose.yml 可以看到整套环境的服务构成服务镜像端口职责redpanda-1redpandadata/redpanda:v24.2.189092Kafka 协议、8082HTTP Proxy、29092容器内通信Kafka 协议兼容的消息服务器jobmanagerpyflink:1.16.0本地构建8081Flink Web UI 与作业提交入口taskmanagerpyflink:1.16.06121、6122容器间 RPC执行算子计算postgrespostgres:145432结果落地存储关键点说明双地址监听Redpanda 同时暴露PLAINTEXT://redpanda-1:29092容器内供 Flink 作业连接与OUTSIDE://localhost:9092宿主机供本地 Python 生产者连接。这正是作业中生产者连localhost:9092、Flink 作业连redpanda-1:29092两种地址并存的根本原因。Flink 镜像本地构建jobmanager通过build: ./Dockerfile.flink构建且设置了pull_policy: never因此首次启动前必须先完成镜像构建。TaskManager 槽位taskmanager.numberOfTaskSlots: 15、parallelism.default: 3意味着作业可默认以 3 并行度运行最多可用 15 个槽位。2.2 PyFlink 镜像的构建内容Dockerfile.flink 以flink:1.16.0-scala_2.12-java8为基础镜像依次完成安装 Python 3.7.9从源码编译PyFlink 1.16 官方支持 Python 3.6–3.8安装 requirements.txt 中的依赖apache-flink1.16.0、psycopg2-binary2.9.1、requests、kafka-python下载四个关键 Jar 到/opt/flink/lib/flink-json-1.16.0.jarJSON 格式解析flink-sql-connector-kafka-1.16.0.jarKafka 连接器flink-connector-jdbc-1.16.0.jarJDBC Sinkpostgresql-42.2.24.jarPostgreSQL 驱动也就是说PyFlink 表 API 之所以能读写 Kafka 和 Postgres完全依赖这几个 connector 与驱动。任何connector not found类报错都应优先检查它们是否已进入/opt/flink/lib/。2.3 启动后的验证访问http://localhost:8081查看 Flink JobManager Web UI用 DBeaver或其他 SQL 客户端连接 Postgres凭据为参数值UsernamepostgresPasswordpostgresDatabasepostgresHostlocalhostPort5432连接成功后先创建第一张落地表作为后续事件进入 Postgres 的着陆区CREATE TABLE processed_events ( test_data INTEGER, event_timestamp TIMESTAMP )这张表与仓库中 start_job.py 里 JDBC Sink 的 DDL 定义完全对应CREATE TABLE processed_events是Kafka → Flink → Postgres链路中 Postgres 侧的接收目标。三、验证 Kafka 连接作业问题 1在真正发送数据前需要先确认 Python 客户端能连通 Redpanda 暴露的 Kafka 端口。3.1 安装客户端pip install kafka-python是否使用独立虚拟环境由你自行决定但建议与 Flink 镜像内的 Python 环境隔离。3.2 连接验证代码在 Jupyter Notebook 或独立脚本中运行import json import time from kafka import KafkaProducer def json_serializer(data): return json.dumps(data).encode(utf-8) server localhost:9092 producer KafkaProducer( bootstrap_servers[server], value_serializerjson_serializer ) producer.bootstrap_connected()producer.bootstrap_connected()返回True表示生产者已成功与 bootstrap server 建立连接。3.3 与仓库生产者实现对照仓库中的 producer.py 是这段连接代码的完整实战版额外展示了发送逻辑import json import time from kafka import KafkaProducer def json_serializer(data): return json.dumps(data).encode(utf-8) server localhost:9092 producer KafkaProducer( bootstrap_servers[server], value_serializerjson_serializer ) t0 time.time() topic_name test-topic for i in range(10, 1000): message {test_data: i, event_timestamp: time.time() * 1000} producer.send(topic_name, valuemessage) print(fSent: {message}) time.sleep(0.05) producer.flush() t1 time.time() print(ftook {(t1 - t0):.2f} seconds)两个实现都遵循同一个模式定义 JSON 序列化函数 → 指定bootstrap_servers→ 构造KafkaProducer→send消息 →flush确保全部送达。注意作业要求不要把 sleep 写进正式的数据发送代码因为 sleep 会人为拉长耗时影响后面问题 3 的时间统计。四、把出租车数据发送到 Kafka topic作业问题 34.1 数据准备与字段裁剪从green_tripdata_2019-10.csv.gz中仅保留作业指定的 7 个字段lpep_pickup_datetimelpep_dropoff_datetimePULocationIDDOLocationIDpassenger_counttrip_distancetip_amount4.2 发送脚本load_taxi_data.py仓库中的 load_taxi_data.py 给出参考实现发送目标为green-datatopic可依作业要求改为green-tripsimport csv import json from kafka import KafkaProducer def main(): # Create a Kafka producer producer KafkaProducer( bootstrap_serverslocalhost:9092, value_serializerlambda v: json.dumps(v).encode(utf-8) ) csv_file data/green_tripdata_2019-10.csv # change to your CSV file path if needed with open(csv_file, r, newline, encodingutf-8) as file: reader csv.DictReader(file) for row in reader: # Each row will be a dictionary keyed by the CSV headers # Send data to Kafka topic green-data producer.send(green-data, valuerow) # Make sure any remaining messages are delivered producer.flush() producer.close() if __name__ __main__: main()实现要点使用csv.DictReader按表头把每一行解析成字典字段名天然与 CSV 列头一致无需手工映射producer.send是异步发送最后的flush()与close()保证缓冲区内消息全部投递后才退出作业要求统计不含 sleep的纯发送耗时可在发送前后记录时间戳相减四舍五入到整数秒。4.3 如何核对数据是否到位在 Flink 作业消费前可以用 Flink Web UIlocalhost:8081观察 topic 的数据量或直接启动下面的消费作业来验证。五、PyFlink 作业的骨架Kafka 源与 JDBC 落地在动手写会话窗口作业之前先理解仓库中已有的 PyFlink 作业结构它是所有后续作业的模板。5.1 表 API 的标准三段式所有 PyFlink 作业start_job.py、aggregation_job.py、taxi_job.py都遵循同一骨架from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import EnvironmentSettings, StreamTableEnvironment # 1. 创建流执行环境并开启 checkpoint env StreamExecutionEnvironment.get_execution_environment() env.enable_checkpointing(10 * 1000) # 2. 创建表环境流模式 settings EnvironmentSettings.new_instance().in_streaming_mode().build() t_env StreamTableEnvironment.create(env, environment_settingssettings) # 3. 执行 DDL 建表 执行 SQLenv.enable_checkpointing(10 * 1000)开启每 10 秒一次的检查点是 Flink 提供精确一次Exactly-Once或至少一次At-Least-Once容错语义的基础生产环境务必保留。5.2 Kafka 源表 DDL以 aggregation_job.py 的 Kafka 源为例def create_events_source_kafka(t_env): table_name events source_ddl f CREATE TABLE {table_name} ( test_data INTEGER, event_timestamp BIGINT, event_watermark AS TO_TIMESTAMP_LTZ(event_timestamp, 3), WATERMARK for event_watermark as event_watermark - INTERVAL 1 SECOND ) WITH ( connector kafka, properties.bootstrap.servers redpanda-1:29092, topic test-topic, scan.startup.mode earliest-offset, properties.auto.offset.reset earliest, format json ); t_env.execute_sql(source_ddl) return table_nameDDL 参数逐个拆解参数含义本作业推荐值connector连接器类型kafkaproperties.bootstrap.serversKafka broker 地址容器内地址redpanda-1:29092topic消费的 topicgreen-tripsscan.startup.mode从何处开始消费earliest-offset从最早消息开始properties.auto.offset.reset无提交 offset 时的重置策略earliestformat消息格式json5.3 JDBC Sink 表 DDL以 start_job.py 的 Postgres Sink 为例def create_processed_events_sink_postgres(t_env): table_name processed_events sink_ddl f CREATE TABLE {table_name} ( test_data INTEGER, event_timestamp TIMESTAMP ) WITH ( connector jdbc, url jdbc:postgresql://postgres:5432/postgres, table-name {table_name}, username postgres, password postgres, driver org.postgresql.Driver ); t_env.execute_sql(sink_ddl) return table_name注意 JDBC URL 使用的是postgres容器服务名而非localhost因为作业运行在 Flink 容器内部必须通过 Docker 网络中的服务名访问 Postgres 容器。5.4 作业提交方式环境就绪后在宿主机上用 Makefile 封装好的命令提交作业# 使用 Makefile推荐 make job # 等价于手动执行 docker-compose exec jobmanager ./bin/flink run -py /opt/src/job/start_job.py --pyFiles /opt/src -d提交聚合作业同理参考 Makefilemake aggregation_job # 等价于 # docker-compose exec jobmanager ./bin/flink run -py /opt/src/job/aggregation_job.py --pyFiles /opt/src -d-d表示后台detached运行提交成功后会在终端看到Job has been submitted with JobID id随后可在 Flink Web UI 的 Running Jobs 页面观察作业状态。六、构建会话窗口作业作业问题 46.1 会话窗口的概念会话窗口Session Window是 Flink 三种窗口Tumbling 滚动、Sliding 滑动、Session 会话之一特点是没有固定长度而是由一个gap间隔参数驱动一旦超过 gap 时长没有新事件到达当前会话即告结束。这天然适合连续运营时段类分析——出租车持续有单算一个活跃会话中间断单超过 gap 就开启新会话。作业要求复制aggregation_job.py并重命名为session_job.py从green-tripstopic 读取数据并修正 schema使用5 分钟 gap 的会话窗口使用lpep_dropoff_datetime作为事件时间水位线容忍 5 秒乱序找出连续无中断运营最长的上下车地点组合。6.2 完整的 session_job.py 参考实现仓库的 2026 届 解决方案 给出了可直接运行的完整实现稍作调整把事件时间改为lpep_dropoff_datetime、聚合键包含DOLocationID即可匹配本作业from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import EnvironmentSettings, StreamTableEnvironment def create_source(t_env): table_name green_trips_session source_ddl f CREATE TABLE {table_name} ( lpep_pickup_datetime VARCHAR, lpep_dropoff_datetime VARCHAR, PULocationID INTEGER, DOLocationID INTEGER, passenger_count INTEGER, trip_distance DOUBLE, tip_amount DOUBLE, event_timestamp AS TO_TIMESTAMP(lpep_dropoff_datetime, yyyy-MM-dd HH:mm:ss), WATERMARK FOR event_timestamp AS event_timestamp - INTERVAL 5 SECOND ) WITH ( connector kafka, properties.bootstrap.servers redpanda-1:29092, topic green-trips, scan.startup.mode earliest-offset, properties.auto.offset.reset earliest, format json ); t_env.execute_sql(source_ddl) return table_name def create_sink(t_env): table_name session_pickup_counts sink_ddl f CREATE TABLE {table_name} ( session_start TIMESTAMP(3), session_end TIMESTAMP(3), PULocationID INT, DOLocationID INT, num_trips BIGINT, PRIMARY KEY (session_start, PULocationID, DOLocationID) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:postgresql://postgres:5432/postgres, table-name {table_name}, username postgres, password postgres, driver org.postgresql.Driver ); t_env.execute_sql(sink_ddl) return table_name def main(): env StreamExecutionEnvironment.get_execution_environment() env.enable_checkpointing(10 * 1000) env.set_parallelism(1) settings EnvironmentSettings.new_instance().in_streaming_mode().build() t_env StreamTableEnvironment.create(env, environment_settingssettings) source_table create_source(t_env) sink_table create_sink(t_env) t_env.execute_sql(f INSERT INTO {sink_table} SELECT window_start AS session_start, window_end AS session_end, PULocationID, DOLocationID, COUNT(*) AS num_trips FROM TABLE( SESSION(TABLE {source_table} PARTITION BY PULocationID, DOLocationID, DESCRIPTOR(event_timestamp), INTERVAL 5 MINUTE) ) GROUP BY window_start, window_end, PULocationID, DOLocationID; ).wait() if __name__ __main__: main()6.3 关键实现点逐项解读1事件时间与水位的声明event_timestamp AS TO_TIMESTAMP(lpep_dropoff_datetime, yyyy-MM-dd HH:mm:ss)是一个计算列把 CSV 中的字符串时间lpep_dropoff_datetime按指定 pattern 解析为TIMESTAMP(3)作为流的事件时间。格式 pattern 必须与 CSV 中实际字符串格式严格一致该数据的典型格式即yyyy-MM-dd HH:mm:ss。2水位线容忍乱序WATERMARK FOR event_timestamp AS event_timestamp - INTERVAL 5 SECOND这行声明允许事件时间最多 5 秒的乱序到达是作业5 秒 tolerance的直接落点。水位线到达后比它更早的迟到事件将被视为迟到数据不参与窗口计算。3会话窗口语法SESSION(TABLE {source_table} PARTITION BY PULocationID, DOLocationID, DESCRIPTOR(event_timestamp), INTERVAL 5 MINUTE)SESSION是会话窗口表函数PARTITION BY指定会话的键——这里按PULocationID, DOLocationID分组保证每个上下车地点组合各自独立成会话DESCRIPTOR(event_timestamp)指定窗口基于的事件时间列INTERVAL 5 MINUTE即 5 分钟 gap某组合超过 5 分钟无新行程会话关闭。4窗口结果字段window_start/window_end是会话窗口的起止时间COUNT(*)统计该会话内包含的行程数。会话越长说明该上下车地点组合连续运营无 5 分钟以上断单的时间越久这正是作业要回答的问题。6.4 查询结果作业运行结束后在 Postgres 中按行程数降序查询即可定位连续运营最久的组合SELECT PULocationID, DOLocationID, num_trips, session_start, session_end FROM session_pickup_counts ORDER BY num_trips DESC LIMIT 3;num_trips最大者对应的(PULocationID, DOLocationID)即为答案。注意提交前先在 Postgres 中创建session_pickup_counts结果表DDL 与 Sink 表定义保持一致否则 JDBC Sink 会因目标表缺失而报错若 Flink 容器内使用 2026 届实现需要把其中的bootstrap.servers改为本环境的redpanda-1:29092、topic 改为green-trips见 docker-compose.yml 中的容器网络配置会话窗口的计算结果只有当对应会话关闭gap 超时触发后才会完整产出因此请等待作业运行一段时间后再查询。七、排错与常用运维命令7.1 常见问题定位现象排查方向生产者bootstrap_connected()返回False确认 Redpanda 容器9092端口已映射、宿主localhost:9092可访问作业报Connector kafka not found检查flink-sql-connector-kafka-1.16.0.jar是否在/opt/flink/lib/作业报 JDBC 连接失败检查 URL 是否使用服务名postgres而非localhost确认 Postgres 容器已就绪窗口迟迟不输出结果检查scan.startup.mode是否为earliest-offset会话窗口要等 gap 超时才会关闭7.2 Makefile 运维命令Makefile 提供了完整的生命周期管理常用目标如下make up # 构建镜像并启动 Flink 集群等价于 docker compose up --build --remove-orphans -d make job # 提交 start_job.py 作业 make aggregation_job # 提交聚合作业 make stop # 停止所有服务 make start # 重新启动所有服务 make down # 停止并移除容器 make clean # 清理容器与悬空镜像7.3 环境清理作业完成后按需清理仓库为只读这些命令作用于你的本地 Docker 环境不影响仓库文件make downPostgres 数据目录若已挂载到本地参考 README.md 中的说明数据会跨容器重启保留需要重置时再手动清理对应数据目录。八、总结本作业完整串起了一条真实可运行的流处理链路Redpanda 负责消息接入Python 生产者写入出租车行程数据PyFlink 以表 API 消费 Kafka、基于事件时间做会话窗口聚合最终把每个连续运营时段的统计结果写入 Postgres。三个关键技术点值得沉淀复用双地址网络模型宿主机生产者连localhost:9092容器内 Flink 作业连redpanda-1:29092这是 Docker 化 Kafka 应用的通用约定事件时间 水位线字符串时间要用计算列解析为时间戳并用WATERMARK ... INTERVAL声明乱序容忍度窗口计算才能正确、完整会话窗口三要素PARTITION BY分组键、DESCRIPTOR事件时间列、INTERVALgap 长度三者共同决定会话的切分逻辑。以此为基础你可以进一步扩展把tip_amount、trip_distance纳入聚合维度做营收会话分析或把结果 Sink 从 Postgres 换成其他 JDBC 目标从而把本实验改造成生产级的流式指标管道。【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表