ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Flink 流式管道实战指南:基于 data-engineer-handbook 的 PyFlink + Kafka + PostgreSQL 全链路搭建

Apache Flink 流式管道实战指南:基于 data-engineer-handbook 的 PyFlink + Kafka + PostgreSQL 全链路搭建 数据工程文档教程【免费下载链接】data-engineer-handbookThis is a repo with links to everything youd ever want to learn about data engineering项目地址https://gitcode.com/GitHub_Trending/da/data-engineer-handbook点击查看免费下载本篇技术指南以># On Ubuntu or Debian: sudo apt-get update sudo apt-get install build-essential # On CentOS or Fedora: sudo dnf install make # On macOS: xcode-select --install # On Windows: choco install make # uses Chocolatey如果不想安装 Make也完全可以直接复制 Makefile 中对应 target 后的命令在终端手动执行——本文后续每个步骤都会同时给出两种方式。获取代码并进入模块目录git clone https://gitcode.com/GitHub_Trending/da/data-engineer-handbook.git cd intermediate-bootcamp/materials/4-apache-flink-training配置凭据从 example.env 到 flink-env.env复制环境文件模块通过环境变量注入 Kafka 凭据与数据库连接信息。第一步是把模板文件复制为实际使用的环境文件cp example.env flink-env.env然后用vim或任意编辑器修改flink-env.envvim flink-env.env环境变量详解example.env 中的完整配置如下KAFKA_WEB_TRAFFIC_SECRETGET FROM WEBSITE KAFKA_WEB_TRAFFIC_KEYGET FROM WEBSITE IP_CODING_KEYMAKE AN ACCOUNT AT https://www.ip2location.io/ TO GET KEY KAFKA_GROUPweb-events KAFKA_TOPICbootcamp-events-prod KAFKA_URLpkc-rgm37.us-west-2.aws.confluent.cloud:9092 FLINK_VERSION1.16.0 PYTHON_VERSION3.7.9 POSTGRES_URLjdbc:postgresql://host.docker.internal:5432/postgres JDBC_BASE_URLjdbc:postgresql://host.docker.internal:5432 POSTGRES_USERpostgres POSTGRES_PASSWORDpostgres POSTGRES_DBpostgres各变量在管道中的实际作用变量说明消费方源码依据KAFKA_WEB_TRAFFIC_KEY/KAFKA_WEB_TRAFFIC_SECRETConfluent Cloud Kafka 的 SASL 认证凭据start_job.py 用于构造PlainLoginModule的 JAAS 配置IP_CODING_KEYip2location.io 地理定位 API 密钥start_job.py 中GetLocationUDF 调用 API 时传入KAFKA_GROUPKafka 消费者组 ID源表/汇表 DDL 中的properties.group.idKAFKA_TOPIC上游事件主题名create_events_source_kafka 读取该主题KAFKA_URLKafka bootstrap servers 地址所有 Kafka 连接器的properties.bootstrap.serversPOSTGRES_URL/JDBC_BASE_URLJDBC 连接串指向宿主机上的 PostgreSQLJDBC sink 的url配置POSTGRES_USER/POSTGRES_PASSWORD/POSTGRES_DBPostgreSQL 连接凭据同时被容器环境变量与 JDBC sink 使用安全警告flink-env.env中保存的是云上 Kafka 资源的真实凭据严禁将其推送或分享到训练营之外否则可能导致云端资源被他人滥用。其余关于凭据的旧版说明可以忽略——仓库更新后需要的一切都已包含在example.env中。如果修改了容器化 PostgreSQL 的POSTGRES_USER与POSTGRES_PASSWORD请保持环境文件与 docker-compose.yml 中的默认值一致否则连接会失败不修改则保持postgres/postgres默认值即可。理解 Flink 镜像构建Dockerfile 逐层拆解在运行管道前先理解 Dockerfile 的构建逻辑这决定了集群的能力边界基础镜像基于flink:1.16.2注意 README 环境变量中标注的FLINK_VERSION1.16.0与镜像实际使用的 1.16.2 存在小版本差异以镜像构建产物为准安装 Python 3.7.9官方 PyFlink 当时仅正式支持 Python 3.6/3.7/3.8而 Debian 11 自带 Python 3.9因此需要从源码编译 3.7.9步骤下载源码 →./configure --enable-shared→make→make install→ 软链python安装 PyFlink 依赖通过 requirements.txt 安装apache-flink1.16.2、psycopg2-binary2.9.1供脚本内 Python 侧使用、requests供 UDF 调用外部 API安装 Java 11并设置JAVA_HOME下载连接器 JAR到/opt/flink/lib/flink-python-1.16.2.jarPyFlink 运行所需flink-sql-connector-kafka-1.16.2.jarKafka 连接器flink-connector-jdbc-1.16.2.jarJDBC 连接器postgresql-42.2.26.jarPostgreSQL JDBC 驱动。这四类 JAR 是后续CREATE TABLE ... WITH (connector kafka / jdbc)能否执行的物质基础任何缺失都会导致作业在运行时抛出 connector 找不到的异常。集群拓扑docker-compose 中的两个核心服务docker-compose.yml 定义了 Flink 会话集群的最小拓扑jobmanager暴露8081:8081Flink Web UI命令为jobmanager通过FLINK_PROPERTIES设置jobmanager.rpc.address: jobmanager使用extra_hosts: host.docker.internal:host-gateway使容器能访问宿主机上的 PostgreSQLtaskmanager依赖 jobmanager 启动命令带--taskmanager.registration.timeout 5 min设置taskmanager.numberOfTaskSlots: 15与parallelism.default: 3为多并行度聚合作业预留资源。两个服务共用image: eczachly-pyflinkpull_policy: never必须由本地build生成并共享卷挂载./src/:/opt/src使作业脚本能被 JobManager 读取。PostgreSQL 不在本 compose 文件中需要单独通过 Makefile 的db-init或训练营前几周的容器启动。运行管道从构建到验证的完整流程第 1 步启动 Flink 集群make up # 没有 make 时手动执行 docker compose --env-file flink-env.env up --build --remove-orphans -d该命令会构建基础镜像并启动 Flink 集群注意make up本身并不包含 PostgreSQL——PostgreSQL 需按前几周教程另行启动或用make db-init。首次构建镜像需要 5 到 30 分钟后续重建只需几秒只要没有删除镜像。务必等待 Flink Web UI 就绪访问 http://localhost:8081/再进入下一步。镜像构建完成后Docker 会自动拉起 jobmanager 与 taskmanager 服务大约需要一分钟。观察容器日志当出现以下日志行时说明 TaskManager 已成功注册到 JobManagertaskmanager Successful registration at resource manager akka.tcp://flinkjobmanager:6123/user/rpc/resourcemanager_* under registration id id_number第 2 步初始化 PostgreSQL 目标表在本地或容器化PostgreSQL 上执行 sql/init.sql创建下游 sink 表CREATE TABLE IF NOT EXISTS processed_events ( ip VARCHAR, event_timestamp TIMESTAMP(3), referrer VARCHAR, host VARCHAR, url VARCHAR, geodata VARCHAR );该表结构与 PyFlink 作业中create_processed_events_sink_postgres定义的 JDBC sink 表字段一一对应是作业能成功写入的前提。第 3 步提交 PyFlink 作业make job # 没有 make 时手动执行 docker compose exec jobmanager ./bin/flink run -py /opt/src/job/start_job.py --pyFiles /opt/src -d大约一分钟后终端会提示作业提交成功例如Job has been submitted with JobID job_id_number。回到 Flink Web UI 的 http://localhost:8081/#/job/running 页面即可看到作业正在运行。第 4 步触发事件并验证落库访问训练营提供的事件触发页面即可向上游 Kafka 主题产生一条新的 Web 事件。随后查询 PostgreSQL 确认数据已写入make psql其实际执行的是进入容器并连接数据库docker exec -it eczachly-flink-postgres psql -U postgres -d postgres在 psql 中验证postgres# SELECT COUNT(*) FROM processed_events; count ------- 739 (1 row)只要计数在增长就证明整条 Kafka → Flink → PostgreSQL 管道已打通。第 5 步停止与清理make stop # 停止正在运行的 compose 服务 make down # 停止并移除 compose 服务 make clean # 移除容器与 none 悬空镜像数据持久化说明PostgreSQL 容器内的/var/lib/postgresql/data挂载到了本机./postgres-data目录因此即使停止或移除容器容器内写入的数据也不会丢失。Makefile 全量命令速查在模块目录运行make help可随时查看所有可用命令当前支持的命令如下命令作用help显示帮助db-init构建并运行 PostgreSQL 数据库服务build构建带 PyFlink 与连接器的 Flink 基础镜像up构建基础镜像并启动 Flink 集群down关闭 Flink 集群job提交 Flink 作业运行 start_job.pyaggregation_job提交聚合作业运行 aggregation_job.pystop/start停止 / 启动 compose 中的所有服务clean停止并移除容器以及 tag 为none的镜像psql在容器内执行 psql 查询 PostgreSQLpostgres-die-mac/postgres-die-pc删除本机Mac / PC与 Docker 中挂载的 postgres 数据目录源码深潜start_job.py 的流式管道实现src/job/start_job.py 是本模块的核心作业其执行流程如下1. 环境初始化与检查点env StreamExecutionEnvironment.get_execution_environment() env.enable_checkpointing(10 * 1000) # 每 10 秒一次检查点 env.set_parallelism(1) settings EnvironmentSettings.new_instance().in_streaming_mode().build() t_env StreamTableEnvironment.create(env, environment_settingssettings)作业以流式模式运行开启每 10 秒的检查点以提供故障恢复能力并行度设为 1。2. 注册自定义 UDFclass GetLocation(ScalarFunction): def eval(self, ip_address): response requests.get(https://api.ip2location.io, params{ ip: ip_address, key: os.environ.get(IP_CODING_KEY) }) 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_typeDataTypes.STRING()) t_env.create_temporary_function(get_location, get_location)GetLocation继承ScalarFunction对每条记录的 IP 发起 HTTP 请求返回国家、州、城市的 JSON 字符串请求失败时返回空对象{}避免单条坏数据中断整个作业。调用失败时返回空 dict 的兜底逻辑response.status_code ! 200体现了流式作业对上游异常的可恢复性设计。3. 声明 Kafka 源表create_events_source_kafka 通过 Flink SQL DDL 声明源表关键配置包括connector kafkatopic与properties.group.id取自环境变量安全协议SASL_SSLPLAIN机制 JAAS 配置使用KAFKA_WEB_TRAFFIC_KEY/SECRETscan.startup.mode latest-offset与properties.auto.offset.reset latest只消费作业启动后的新事件计算列event_timestamp AS TO_TIMESTAMP(event_time, yyyy-MM-ddTHH:mm:ss.SSSZ)将字符串事件时间解析为时间戳format json。4. 声明 PostgreSQL 汇表create_processed_events_sink_postgres 声明 JDBC sinkconnector jdbc, url os.environ.get(POSTGRES_URL), table-name processed_events, username / password 从环境变量读取 driver org.postgresql.Driver注意url使用host.docker.internal——这正是 docker-compose.yml 中extra_hosts映射宿主机网关的用武之地使容器内作业可以访问宿主机上的 PostgreSQL。5. 组装 INSERT 查询t_env.execute_sql(f INSERT INTO {postgres_sink} SELECT ip, event_timestamp, referrer, host, url, get_location(ip) as geodata FROM {source_table} ).wait()每条 Kafka 事件经get_location(ip)增强后写入processed_events.wait()阻塞至作业完成提交。进阶扩展aggregation_job.py 的窗口聚合src/job/aggregation_job.py 展示了在流上做滚动窗口Tumbling Window聚合的写法源表通过计算列window_timestamp AS TO_TIMESTAMP(event_time, ...)解析事件时间并用WATERMARK FOR window_timestamp AS window_timestamp - INTERVAL 15 SECOND声明 15 秒的水位线容忍乱序数据作业以并行度 3 运行env.set_parallelism(3)对应 taskmanager 的parallelism.default: 3对每个 5 分钟窗口按host分组统计num_hitsTumble.over(lit(5).minutes).on(col(window_timestamp))同时按hostreferrer双维度分组写入第二张聚合表processed_events_aggregated_source结果通过 JDBC 写入 PostgreSQL。make aggregation_job即可提交该作业。这一示例揭示了从「原始事件管道」到「指标聚合管道」的演进路径也是本模块作业的核心素材——按 IP 与 host 做 5 分钟 gap 的会话化sessionization并回答「Tech Creator 上单个用户会话的平均事件数」等问题。验证与故障排查要点UI 未就绪检查docker compose logs jobmanager确认 TaskManager 注册日志出现后再提交作业Kafka 认证失败核对flink-env.env中KAFKA_WEB_TRAFFIC_KEY/SECRET与 JAAS 格式注意转义引号写入失败确认已执行 sql/init.sql 且 PostgreSQL 凭据与 compose 环境变量一致无数据消费确认上游事件已触发且scan.startup.modelatest-offset下作业需在事件产生前启动作业崩溃于外部 APIGetLocation的非 200 兜底可防止单条异常记录阻塞管道排查时可先观察 Flink UI 的异常栈与检查点状态。至此你已经掌握了基于本仓库 Flink 训练模块的完整流式管道从环境准备、凭据配置、镜像构建到作业提交、窗口聚合与结果验证。这套模式可直接迁移到你自己的 Kafka Flink PostgreSQL 实时数据处理场景中。赞分享数据工程文档教程【免费下载链接】data-engineer-handbookThis is a repo with links to everything youd ever want to learn about data engineering项目地址https://gitcode.com/GitHub_Trending/da/data-engineer-handbook点击查看免费下载相关推荐Data Engineering Zoomcamp用 Apache Flink 构建端到端 PyFlink 流式管道实战指南Data Engineering Zoomcamp用 Apache Flink 构建端到端 PyFlink 流式管道实战指南 Apache Flink 是当前教程数据工程GB28181开源视频监控平台多品牌摄像头一套平台搞定免费开箱即用十分钟能看到画面吗GB28181开源视频监控平台多品牌摄像头一套平台搞定免费开箱即用十分钟能看到画面吗 wvp GB28181 pro 是一款免费开源、可商用的 GB28后端音视频前端Apache Flink PyFlink 指南用 Python API 构建批流一体的数据管道Apache Flink PyFlink 指南用 Python API 构建批流一体的数据管道 PyFlink 是 Apache Flink 的 Python后端大数据流处理批处理上一篇Android Studio 中文语言包安装全记录不写一行代码的全界面汉化方案下一篇免费商用中文字体怎么选思源宋体CN 7个字重从下载到网页上线的完整实操创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表