ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

PyFlink Table API 纯 Python 实现 CSV WordCount 全流程拆解

2026/9/8 2:36:42 拓冰建站 浏览量
PyFlink Table API 纯 Python 实现 CSV WordCount 全流程拆解 我最早看到“PyFlink Table API 用纯 Python 写一个 WordCount读 CSV 聚合 写出”这个标题的时候第一反应是又一个入门 Demo。但真按这个链路在本地跑通之后我发现它其实把 PyFlink 最核心的几条主线都串起来了——环境怎么搭、TableEnvironment 怎么用、SQL 和 Table API 怎么结合、文件连接器怎么配、批流模式有什么差别。更关键的是整个过程里我一行 Java 都没写业务逻辑全是 Python这对我这种从 Python 生态转过来的开发者来说价值比想象中大得多。这篇文章就围绕这个 WordCount 案例展开适合三类人看一是刚接触 PyFlink、想知道 Python 能不能正经跑 Flink 作业的二是已经会 DataStream 但被 Java 折腾得够呛、想试试 Table API 的三是纯粹想用 Python 对 CSV 做快速清洗、聚合、落盘的数据开发。我会把环境准备、DDL 建表、聚合逻辑、结果写出、常见报错全部拆开讲重点解释每一步为什么这么做最后再分享几个真正踩过坑之后才知道的细节。1. 先想清楚这里说的“纯 Python”到底是什么意思很多人第一次看到“用纯 Python 写 Flink”会有两个极端反应要么觉得是玩具要么觉得以后不用碰 Java 了。这两种理解都不太准确。要搞清楚 PyFlink Table API 的定位得先把 WordCount 为什么是入门第一课、Table API 和 DataStream API 到底差在哪、以及整个数据链路对应的计算模型讲明白。1.1 WordCount 为什么是入门第一课数据领域的“Hello World”就是 WordCount它看起来简单但数据处理的三个核心环节全都有读取数据源、做计算转换、把结果写出去。放到 Flink 里就是 Source、Transform、Sink 三段式结构。在 PyFlink 之前想写一个 Flink 的 WordCount一般得用 Java 或 Scala写一个main方法、配好 Maven 依赖、打包、提交到集群。光这一套流程就劝退了不少只熟悉 Python 的人。而用 PyFlink Table API 写 WordCount整个诉求非常朴素从一个 CSV 文件里把内容读进来按单词分组统计再把统计结果写到另一个 CSV 文件。这个案例好在哪第一它足够短核心逻辑只有一张建表 DDL 加一句聚合 SQL第二它完整覆盖了数据文件从进到出的全链路跑通它你对 PyFlink 的作业生命周期就有概念了第三它很容易验证结果打开输出文件扫一眼就知道对不对不需要依赖可视化界面。所以我一直觉得WordCount 不是“幼稚的 Demo”而是用最小成本验证整套工具链能不能用的试金石。1.2 Table API 和 DataStream API 怎么选PyFlink 其实同时提供了两套 APIDataStream API 和 Table API/SQL。最开始做 PyFlink 的时候DataStream API 是主力但它在 Python 场景下写起来非常啰嗦。你要处理一个时间窗口、要管理状态、要做 keyBy每一步都是显式的操作代码量大不说调试也麻烦。更现实的问题是PyFlink 的 DataStream API 在 Python UDF 和底层算子之间来回切换性能损耗和心智负担都不小。Table API 的出现就是为了解决这个问题。它把数据抽象成一张“表”把计算抽象成“查表”。你不需要关心数据是怎么分区的、算子是怎么串联的只需要描述希望得到什么结果。比如这次我们要做的分组计数本质就是一条SELECT word, COUNT(1) FROM source GROUP BY word。这里有一个特别重要的点Table API 和 SQL 是统一的Table API 表达不清晰的场景可以直接写 SQLSQL 表达不了的自定义逻辑可以注册成 UDF 混着用。而且 Table API 是批流统一的同一套查询逻辑在流模式下是持续计算在批模式下是一次性执行API 层面的写法几乎不变。选型建议也比较直接如果你要做的是数据清洗、聚合、关联、报表统计这类声明式任务优先考虑 Table API如果你要做的是一些平台底层算子优化、自定义状态管理、精细控制事件时间处理再考虑 DataStream。实际生产里绝大多数 ETL 任务用 Table API 已经足够而且代码量能减少一半以上。1.3 核心链路读 CSV、聚合、写出分别是什么整个 WordCount 其实是一条清晰的数据管道。第一步“读 CSV”对应的是 Source数据从外部文件系统进入 Flink这个过程需要告诉 Flink 三件事文件在哪、文件是什么格式、文件里每一列是什么类型。第二步“聚合”对应的是 Transform按word字段分组然后对组内记录做计数这个动作会把相同 key 的数据汇聚到同一个逻辑分组里。第三步“写出”对应的是 Sink把计算结果写回文件系统同样需要声明目标路径、格式、字段结构。在 Table API 里这三个环节被抽象成“源表”“查询”“结果表”。源表是通过CREATE TABLE语句声明的虚拟表查询是对虚拟表执行 SQL 或 Table API 操作结果表定义了最后写入的位置。这个抽象非常接近普通数据库的“表”概念所以对经常用 Pandas、SQL 的人来说上手 PyFlink Table API 的曲线其实并不陡。有一点需要明确Table API 的“表”和关系型数据库的表并不完全一样底层仍然是分布式流处理和批处理引擎只是 API 层帮你把复杂性包装起来了。你不需要手动去map、flatMap但引擎依然会对数据进行分区、并行计算、shuffle。理解了这一点后面遇到“为什么输出有多个文件”“为什么并行度会影响结果文件数”这类问题就不会觉得奇怪了。2. 环境准备本地跑起来只需要三步很多入门教程会直接跳过环境准备但 PyFlink 的环境坑真的不少。我见过有人在安装环节卡了一整天的也有明明装好了却因为 Java 版本不对一直报错的。所以这里我把环境部分单独拿出来讲确保你照着做能一次跑通。2.1 版本选择与安装PyFlink 的安装很简单用 pip 就行但版本选择有一点讲究。PyFlink 版本和 Flink 版本是绑定的比如apache-flink1.18.1对应的就是 Flink 1.18.1。安装方式pip install apache-flink1.18.1建议在虚拟环境里装避免把系统 Python 环境搞乱。我自己习惯用 conda 建一个独立环境conda create -n pyflink-demo python3.10 -y conda activate pyflink-demo pip install apache-flink1.18.1Python 版本要特别注意PyFlink 对 Python 版本有要求一般是 3.8 到 3.11 之间。装太新的 Python 版本可能导致找不到对应的 PyFlink 依赖包装太旧的又可能跟其他库冲突。1.18 用 Python 3.10 是比较稳的组合。版本选择上我的建议是“选新不选旧但别追最新”。太老的版本比如 1.14、1.15 有一些 API 已经变了网上教程你可能对不上最新版本虽然功能多但生态里的第三方库适配可能滞后。1.17 或 1.18 目前看性价比最高生产环境也在大面积使用。2.2 本地运行还需要 Java这里有坑这是“纯 Python”最容易产生误解的地方。业务代码确实全用 Python 写但 PyFlink 底层运行的仍然是 JVM 上的 Flink 引擎所以你的机器上必须得有 Java。如果没装 JDK运行时会直接报java.lang.NoClassDefFoundError或者提示找不到 Java。JDK 版本建议 8 或 11PyFlink 1.18 用 Java 11 没毛病。安装完之后配置环境变量# Linux / macOS export JAVA_HOME/path/to/jdk export PATH$JAVA_HOME/bin:$PATH # Windows 在系统环境变量里加 JAVA_HOME并新建 %JAVA_HOME%\bin 到 PATH很多人在这一步会犯一个错只装了 JRE 没装 JDK。Flink 运行调试时需要javac相关的类库只装 JRE 有时会出莫名其妙的问题。所以直接装完整 JDK 最省事。还要有个心理准备第一次跑 PyFlink 作业时启动 JVM、加载依赖 jar、初始化执行环境整个过程可能需要十几秒甚至更久。这不是卡死了也不是死循环只是 Flink 在做初始化。刚开始跑通之前不要频繁中断否则容易误判是程序问题。2.3 验证环境是否可用在正式写代码之前先做两步快速验证。第一步查看 PyFlink 能不能正常导入python -c from pyflink.table import EnvironmentSettings, TableEnvironment; print(ok)如果这行能输出ok说明 PyFlink 安装成功。第二步准备一个测试用的 CSV 文件。我习惯在项目目录下建一个data文件夹放一个很简单的input.csv内容如下flink python flink table api python flink csv后续所有代码都基于这个文件来跑。文件准备好之后我们可以开始写核心逻辑了。这一步很关键我建议你别直接用后面那种好几列数据的大文件来测试先确保最简链路能跑通再逐步加复杂度。3. 从 CSV 到聚合结果核心代码逐段拆解现在到了文章最核心的部分。我会把代码拆成四段讲创建 TableEnvironment、定义 CSV 源表、写聚合查询、定义写出目标表。最后给出一份可以直接跑的完整脚本。每一段我都会解释关键参数的意义而不是让读者直接抄完就完事。3.1 第一步创建 TableEnvironmentTableEnvironment 是 PyFlink Table API 的入口类似于 Spark 里的 SparkSession。所有建表、查询、写出的动作都要通过它来发起。创建方式如下from pyflink.table import EnvironmentSettings, TableEnvironment env_settings EnvironmentSettings.in_streaming_mode() t_env TableEnvironment.create(env_settings)这里用到了in_streaming_mode()也就是流模式。你可能会想我们明明是离线读文件做统计为什么用流模式这是因为 Flink 是流处理出身Table API 默认的底层运行模式就是流式。在流模式下处理有界数据比如一个 CSV 文件本质上是把文件当成一个有边界的流来处理。如果你确定只用批模式也可以改成EnvironmentSettings.in_batch_mode()。批模式在数据倾斜处理、结果产出方式上会有些区别但对这样的小案例两者结果一致。我为什么推荐先用流模式因为 PyFlink 的流模式兼容性最好后面你想接 Kafka、接实时数据源都不用改环境配置。这是 Table API 批流统一带来的红利。还有一个细节TableEnvironment.create()之后不需要手动 start 或 close脚本执行完进程退出即可。如果你在 Jupyter Notebook 里跑每个 cell 之间建议复用同一个t_env不要反复创建否则会重复加载环境很慢。3.2 第二步把 CSV 定义成一张表读取 CSV 在 Flink 里的标准做法是声明一张源表用的是CREATE TABLEDDL。把数据源抽象成表是 Table API 对程序员最大的友好之处。t_env.execute_sql( CREATE TABLE source_table ( word STRING ) WITH ( connector filesystem, path data/input.csv, format csv ) )这段 DDL 里最核心的是WITH部分的三个参数。connector设置成filesystem表示数据来源是文件系统这是 Flink 内置的文件连接器本地文件、HDFS、S3 都可以用它。path就是文件路径这里支持相对路径也可以写绝对路径。format设置为csv表示用 CSV 格式解析文件内容。字段声明那块也值得注意我声明了一个word STRING。这意味着 CSV 文件每行会被解析成一个字段字段名是word类型是字符串。如果 CSV 文件有表头怎么办最简单的方式是建表时把表头字段当成普通字段做映射。更优雅的做法是后续通过参数跳过表头但那个参数不是所有版本都支持为了稳我建议你准备 CSV 时第一行就放数据不要放标题。顺带说一句如果你的 CSV 是标准逗号分隔这个配置就够用了如果是制表符或其他分隔符还需要额外加csv.field-delimiter \t之类参数。我自己用管道符|的场景也遇到过都可以通过这个参数解决。3.3 第三步聚合统计表定义好之后聚合就非常轻松了。最直观的方式是直接写 SQLresult_table t_env.sql_query( SELECT word, COUNT(1) AS cnt FROM source_table GROUP BY word )这里COUNT(1)是统计每个分组有多少行GROUP BY word是按单词分组。很多从 Pandas 转过来的人会特别不适这不就是df.groupby(word).size()吗对本质就是这件事但 Flink 是分布式执行数据量大的时候性能不是一个量级。如果你想用更“API 风格”的写法不用 SQL 其实也可以source_table t_env.from_path(source_table) result_table source_table.group_by(word).select(source_table.word.alias(word), source_table.word.count.alias(cnt))两种写法底层生成的执行计划几乎一样选哪种纯粹看个人习惯。我个人更推荐 SQL 方式因为阅读成本低团队协作时也更容易评审逻辑。还有一个容易困惑的地方COUNT(1)的结果类型是什么在 Flink 里是BIGINT对应 Python 的int。所以后面建结果表的时候计数字段声明的类型也要和它匹配否则写入阶段会报类型不一致的错误。3.4 第四步把结果写回 CSV聚合结果现在是一张虚拟表result_table它还没有真正落盘。要写出到 CSV同样需要声明一张结果表t_env.execute_sql( CREATE TABLE sink_table ( word STRING, cnt BIGINT ) WITH ( connector filesystem, path data/output, format csv ) )注意这里的path我写的是data/output它不是一个具体文件而是一个目录。这是 Flink 文件系统 sink 的设计输出会生成目录下的一个或多个文件。如果你指定的是一个不存在的路径Flink 会帮你创建目录如果目录已经存在不同版本行为可能不同有的会报错有的会覆盖这个点后面我会单独讲。然后执行写出动作result_table.execute_insert(sink_table).wait()execute_insert是把result_table的结果插入到sink_table。这个方法返回一个JobExecutionResult调用.wait()是为了等待作业执行完成。如果不加.wait()脚本可能在作业还没跑完的时候就走到了最后一行然后退出导致输出目录为空或者结果不完整。这是一个非常典型的坑我见不少人栽在这里。3.5 完整脚本把上面几段拼起来就是一份可以直接运行的完整脚本from pyflink.table import EnvironmentSettings, TableEnvironment # 1. 创建执行环境 env_settings EnvironmentSettings.in_streaming_mode() t_env TableEnvironment.create(env_settings) # 2. 定义 CSV 源表 t_env.execute_sql( CREATE TABLE source_table ( word STRING ) WITH ( connector filesystem, path data/input.csv, format csv ) ) # 3. 聚合统计 result_table t_env.sql_query( SELECT word, COUNT(1) AS cnt FROM source_table GROUP BY word ) # 4. 定义 CSV 结果表 t_env.execute_sql( CREATE TABLE sink_table ( word STRING, cnt BIGINT ) WITH ( connector filesystem, path data/output, format csv ) ) # 5. 执行写入并等待完成 result_table.execute_insert(sink_table).wait()这份脚本放到项目根目录下命名为wordcount.py配合前面那个data/input.csv理论上可以直接跑通。下一节我会带你走一遍执行和验证流程包括怎么确认结果、怎么处理“一行一句英文”的进阶场景。4. 跑起来执行、调试与结果校验代码写完了接下来就是见证结果的时刻。这一节我会从执行命令开始到结果文件检查再到一个比较常见的进阶场景如果 CSV 里不是单词而是一个个英文句子怎么做真正的分词统计。4.1 命令行执行与执行参数在项目根目录执行python wordcount.py第一次运行大概率会等一会儿因为 Flink 要初始化 JVM、加载内部依赖。等到脚本正常结束在data/output目录下会生成文件。用ls可以看到类似下面的内容data/ ├── input.csv └── output/ ├── part-xxx-0 ├── part-xxx-1 └── ...part-开头的文件就是结果文件。Flink 文件系统 sink 默认按照并行度生成多个文件如果你只有一个 CPU 核心或者没配并行度通常只有一个part文件如果机器核数多Flink 默认并行度会取可用核心数结果文件也会对应变多。文件数量多不是错误这是分布式计算的正常现象。如果你想强制让输出合并成一个文件可以在创建执行环境后设置并行度t_env.get_config().set(parallelism.default, 1)但这种做法不建议在正式大数据场景用因为会牺牲并行处理能力。本地小数据集为了看着清爽可以这样设。4.2 结果文件长什么样打开任意一个part文件内容应该是类似这样的flink,3 python,2 table,1 api,1 csv,1第一列是单词第二列是出现次数。我用的是逗号分隔因为 CSV format 默认分隔符就是逗号。看到这个结果说明整个读、聚合、写链路已经全部打通。如果你发现结果文件和预期不一致优先检查三件事第一源 CSV 里单词是不是有多余空格或空行第二路径有没有写错Flink 对找不到文件时的报错信息有时候不够直观第三是不是用了流模式但没.wait()导致程序提前退出。这三个问题覆盖了 80% 的首次运行失败场景。4.3 进阶场景CSV 里存的是句子怎么办刚才我们假设 CSV 每行是一个单词这其实是简化版。真正的 WordCount 通常会处理整段文本比如hello world flink table api hello flink每行是一条句子需要先分词再聚合。这个时候只靠内置 CSV 格式和 SQL 就不够了需要写一个 Python UDF 来做分词。PyFlink 支持用 Python 写 UDF而且注册和调用都很简单。完整示例如下from pyflink.table import DataTypes from pyflink.table.udf import udf udf(result_typeDataTypes.ARRAY(DataTypes.STRING())) def split_words(line: str): return [word for word in line.lower().split() if word] t_env.create_temporary_function(split_words, split_words)注册之后SQL 里就可以这么查SELECT word, COUNT(1) AS cnt FROM source_table CROSS JOIN UNNEST(split_words(line)) AS t(word) GROUP BY word这里的UNNEST是 Flink SQL 里把数组展开成多行的函数配合CROSS JOIN就能把一行句子拆成多行单词。这个示例的价值在于它展示了纯 Python 和 Table API 是怎么协同工作的你甚至可以在 UDF 里用正则、用第三方库如jieba做中文分词。这就是 PyFlink 对 Python 开发者最大的吸引力——复杂逻辑用 Python 写分布式的脏活累活交给 Flink 干。5. 常见问题排查与避坑实录跑了这么多遍踩坑经验自然也攒了不少。这一节我把最常见的问题整理成一张速查表再单独说几个常规文档里不怎么会写到的经验。5.1 高频报错速查表报错现象可能原因解决办法ModuleNotFoundError: No module named pyflink没安装 PyFlink 或装到了其他环境确认在虚拟环境里执行pip install apache-flinkjava.lang.ClassNotFoundException或找不到 Java缺少 JDK 或JAVA_HOME没配置安装 JDK 8/11配置JAVA_HOME和PATHFileNotFound/Path does not exist源文件路径不对使用绝对路径或确认相对路径是相对于执行命令的目录结果文件为空没有.wait()或作业没跑完补上.wait()等待作业执行完成写入时目录已存在报错Flink 对输出目录存在性敏感删除已存在的输出目录或换一个输出目录路径SQL 解析失败 / 表找不到DDL 或查询 SQL 有语法问题检查表名、字段名是否一致先execute_sql建表再sql_query类型不匹配BIGINT与INT冲突聚合结果类型和结果表字段类型不一致统一切换成BIGINT或修改 SQL 里的CAST表达式这张表里的问题我基本都遇到过。尤其是输出目录已存在的情况有时候不是每次都会报错而是要看目录里有没有旧文件非常容易让人困惑。最稳妥的解决办法是每次运行前先删掉旧的输出目录。5.2 我踩过的几个坑和经验第一个坑是关于路径的。PyFlink 里的相对路径并不是相对于 Python 脚本所在目录而是相对于你执行python命令时所在的当前工作目录。我第一次写脚本时把data/input.csv写在项目子目录下命令行却跑在根目录结果一直提示找不到文件。排查了很久才意识到是当前工作目录的问题。建议脚本里统一用绝对路径或者在执行脚本前先cd到项目根目录这两个做法都能避免这种问题。第二个坑是环境变量。我有一次在朋友电脑上演示时他的机器有多个版本的 JavaFlink 优先找到了 Java 17结果跑出来的行为很奇怪有些参数不生效。后来把JAVA_HOME明确指到 Java 11 才正常。所以在环境准备阶段一定要确认JAVA_HOME指向的是正确版本而不是“只要有 Java 就行”。第三个经验是关于测试数据的。本地调试时建议先用非常小的文件最好就三四行这样第一能快速跑完第二结果一眼就能看出来对不对。如果一上来就搞一个几万行的 CSV程序跑挂了都不知道是逻辑问题还是数据问题。小数据验证通过之后再逐步加大数据量观察执行时间和内存表现。第四个经验是关于 Python 虚拟环境的。PyFlink 依赖包数量比较多如果直接装在全局环境很容易跟已有的 pandas、numpy 版本冲突。我见过有人装了 PyFlink 之后原来能正常跑的 pandas 代码突然报错就是因为 PyFlink 引入了某个传递性依赖把 pandas 的底层库升级了。虚拟环境虽然不能完全规避这种问题但至少不会影响系统级环境出了问题也能直接删掉重建。最后一个想分享的经验是不要急着上集群。很多人学 PyFlink 第一反应是部署一套 Flink 集群其实在本地单机跑通这些案例已经能覆盖大部分 API 使用场景。集群部署之后带来的网络、权限、资源调度问题对学习初期反而是负资产。先在本地把读 CSV、聚合、写出这条链路跑熟理解作业生命周期和日志排查手段再考虑怎么上生产环境这个节奏会舒服很多。对我个人来说PyFlink Table API 解决了一个很实际的痛点它能让我用 Python 的表达力去写分布式计算的任务。这个 WordCount 虽然简单但它把 Python 生态和 Flink 引擎之间的桥梁完整地搭了起来。项目后续如果想扩展可以尝试用同一套 API 把数据源换成 Kafka、把结果写到 ClickHouse或者用 Catalog 管理多张表。方向很多但起点就是今天这个 CSV 小文件。你自己动手跑一遍感受会比我写再多字都来得直接。