ARTICLE DETAIL

资讯详情

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

Flink/PyFlink读写CSV实战:Schema配置、时间格式与脏数据避坑指南

Flink/PyFlink读写CSV实战:Schema配置、时间格式与脏数据避坑指南 最近在做一个数据接入的活儿客户丢过来一批 CSV 文件几十个字段、三种时间格式、还夹杂着空值和脏行要求用 Flink 和 PyFlink 处理完再落到库里。CSV 这个格式看起来人畜无害可真到了 Flink/PyFlink 里做批式或流式读写时坑基本都藏在 Schema 配置和解析细节里。这篇文章把 Flink/PyFlink 读写 CSV 的常用姿势、Schema 高级配置以及几个容易踩的雷一次性讲清楚适合正在用 Flink SQL、PyFlink 接 CSV 文件数据或者打算把 CSV 同步进 MySQL、ClickHouse 的朋友参考。先说结论Flink 处理 CSV 的能力不止一种但生产环境里我基本只用 Table API / SQL DDL csv format 这一条路。它最大的价值是让你用声明式的方式把“CSV 文本长什么样”完整描述清楚剩下的序列化、反序列化、类型转换、引号转义、容错策略全部交给框架处理。下面把里面的门道一条条拆开讲。1. 先搞清楚Flink / PyFlink 读写 CSV 有哪几条路1.1 CSV Format 不是连接器而是表格格式很多刚接触 Flink 的朋友会把“CSV Format”误当成一个连接器其实它本身不负责读文件、不负责连 Kafka也不负责写数据库。它的角色是 Table Format也就是“表格格式转换器”必须配合 Filesystem、Kafka、JDBC 这类连接器一起使用。连接器负责拿数据和放数据CSV Format 负责把每条数据的文本形态和内部结构互相转换。我举一个最常见的组合connector filesystem负责扫描/data/events/*.csv这批文件读到每一行文本后交给format csv去解析成 Flink 内部的行对象。反过来写入时csv format 再把行对象拼成一行 CSV 字符串交给 filesystem connector 落地。所以你在 DDL 里看到的WITH参数一部分属于连接器比如 path另一部分属于 format比如 csv.field-delimiter两类参数混在一起理解它们的归属是排查问题的起点。底层实现也不是你用String.split(,)写完就完事的那种弱解析flink-csv 模块里的解析器实现了完整的 CSV 语义支持双引号包裹、引号内转义、注释、数组分隔符、null 字面量等规则。这也是为什么它能处理“字段内容里带逗号和换行符”的合法 CSV而你手写 split 遇到这种数据直接错位。1.2 三种姿势对比手写解析、SQL DDL、DataStream我在实际项目里见过三种处理 CSV 的路线路线常见做法优点缺点DataStream 手写解析读文本流自己 split、自己转型、自己处理引号灵活适合一次性脚本所有脏数据规则都要自己写转义、类型异常、空值处理很容易漏Table API / SQL DDL csv formatCREATE TABLE 声明 SchemaINSERT INTO SELECT 完成转换声明式、批流统一、类型自动映射、容错参数开箱即用遇到官方 CSV Format 不支持的场景比如表头需要绕路PyFlink DataStream Datastream API在 Python 里用 map 函数逐行处理适合写 Python 自定义逻辑类型标注麻烦流式计算性能损耗明显代码可维护性差我的选择偏好非常明确凡是能用 SQL DDL 表达的绝对不用手写解析。因为 CSV 的解析规则实在太细你稍微漏一个引号转义或者日期格式线上就会出脏数据。而 SQL DDL 把“每列什么类型、怎么解析、怎么容错”全部显式表达出来别人接手也容易看懂。PyFlink 用户也优先走execute_sql不要绕到 DataStream 里去做。1.3 依赖配置Java 和 PyFlink 各要注意什么Java 工程里用 CSV Format核心依赖是flink-csv。如果你用 Maven加上这么一段dependency groupIdorg.apache.flink/groupId artifactIdflink-csv/artifactId version你的Flink版本/version /dependency如果用的是 Flink Table API还需要flink-table-api-java-bridge这通常和运行环境版本保持一致。PyFlink 的情况稍微特殊一点。官方发布的 PyFlink wheel 包一般会带常用 format 的 jar理论上你不需要额外操作。但如果你是自己精简过的发行包或者客户环境里裁剪过 lib 目录运行时会报Could not find any factory for identifier csv。这个报错基本就是缺 flink-csv.jar把对应版本的 jar 放进${FLINK_HOME}/lib目录基本能解决。2. 五分钟跑通PyFlink 里用 SQL DDL 读 CSV2.1 一张源表的完整 DDL我们直接看一个例子。假设有一批用户事件 CSV路径是/data/events/csv/每行数据长这样10001,zhangsan,88.50,true,click;share,2024-01-01 12:30:00对应 DDL 如下CREATE TABLE user_event_csv ( user_id BIGINT, user_name STRING, score DECIMAL(10, 2), is_vip BOOLEAN, tags ARRAYSTRING, event_time TIMESTAMP(3) ) WITH ( connector filesystem, path /data/events/csv/*.csv, format csv, csv.field-delimiter ,, csv.ignore-parse-errors false, csv.timestamp-format yyyy-MM-dd HH:mm:ss );这里的path支持通配符/data/events/csv/*.csv会匹配目录下所有 CSV 文件这是 Filesystem 连接器的能力。csv.field-delimiter指定列分隔符默认就是逗号但如果你的文件是竖线、Tab 分隔这里改成csv.field-delimiter |或csv.field-delimiter \t就行。csv.timestamp-format对应 event_time 的解析格式因为样例里时间没有毫秒用默认的yyyy-MM-dd HH:mm:ss就够了。有一点必须强调Flink 的 csv format 默认不读表头。它把文件每一行都当作数据如果 CSV 第一行是user_id,user_name,...那么整行会被解析器尝试按字段类型去转换然后要么报错要么解析成 null结果完全不可控。这个问题我在 2.3 节专门讲。2.2 PyFlink 环境初始化与查询写入PyFlink 里跑这个 DDL 非常简单关键代码不超过 20 行from pyflink.table import EnvironmentSettings, TableEnvironment # 创建 Table 环境流模式足够批式场景也能复用 env_settings EnvironmentSettings.in_streaming_mode() t_env TableEnvironment.create(env_settings) # 注册源表 t_env.execute_sql( CREATE TABLE user_event_csv ( user_id BIGINT, user_name STRING, score DECIMAL(10, 2), is_vip BOOLEAN, tags ARRAYSTRING, event_time TIMESTAMP(3) ) WITH ( connector filesystem, path /data/events/csv/*.csv, format csv, csv.field-delimiter ,, csv.ignore-parse-errors false, csv.timestamp-format yyyy-MM-dd HH:mm:ss ) ) # 先做一条简单查询验证解析是否正常 result t_env.sql_query(SELECT user_id, user_name, score FROM user_event_csv) result.execute().print()如果本地跑能看到查询结果说明解析链路是通的。接下来要做得最多的操作是直接INSERT INTO SELECT把 CSV 源写到某个目标端。整条链路不需要写任何 Java 代码这也是 PyFlink 处理这类任务最舒服的地方。2.3 写入 CSVsink 定义与引号规则把处理结果写回 CSV 文件同样用 DDL 声明 sinkCREATE TABLE user_event_csv_sink ( user_id BIGINT, user_name STRING, score DECIMAL(10, 2), event_time TIMESTAMP(3) ) WITH ( connector filesystem, path /data/events/csv_out, format csv );然后执行写入INSERT INTO user_event_csv_sink SELECT user_id, user_name, score, event_time FROM user_event_csv;写入侧有一个细节容易被忽略CSV 写入并不是“字段中间有逗号就自动加引号”Flink 的 writer 会判断字段值里是否包含分隔符、双引号、换行符等特殊字符一旦包含就会自动用双引号包裹并把包裹内容里的双引号通过双引号转义。这是 CSV 协议的标准行为不是 Flink 特有的 bug。自己拼 CSV 字符串时千万别漏掉这一层不然下游工具解析出来的列顺序直接乱掉。2.4 带表头的 CSV 怎么处理这是被问得最多的问题之一。官方的 CSV Format 没有“跳过表头”这样的参数它把每一行都当数据。所以要处理带表头文件我常用的方案有三个第一数据接入前先预处理用 shell 一行tail -n 2或者 Python 脚本去掉表头。这个最粗暴也最可靠。第二如果表头文件不多可以把表头文件和真正的数据文件放到不同目录或者利用 path 通配符只匹配数据文件。第三如果读取的是单一大文件且实在没办法改动可以退到 DataStream 场景用read_text_file把第一行单独丢弃再走 CSV 解析。我特别不建议用csv.ignore-parse-errors true去“跳过”表头。这个参数做的事情是把解析失败字段置成 null并不代表整行被丢弃表头里的字符串字段可能会被当成合法 STRING 处理后续聚合结果全是错的而且排查起来非常隐蔽。3. Schema 高级配置字段类型、时间格式、脏数据全解析3.1 Flink 类型和 CSV 文本的映射关系CSV 是纯文本格式没有任何列类型信息。因此 Schema 就是解析的唯一依据DDL 里你把字段声明成什么类型解析器就按什么类型去转换。这是 CSV Format 使用的核心逻辑。Flink 字段类型CSV 文本示例说明STRINGhello原样保留BOOLEANtrue / false大小写不敏感INT / BIGINT123不能带千分位不能带引号DECIMAL123.45读取时可按普通数字字符串解析DATE2024-01-01默认格式 yyyy-MM-ddTIME12:30:00默认格式 HH:mm:ssTIMESTAMP2024-01-01 12:30:00默认格式 yyyy-MM-dd HH:mm:ssARRAYSTRINGa;b;c默认用分号分隔不是逗号MAP / ROW不建议使用CSV 对这种嵌套结构支持有限上表里最容易踩坑的是ARRAY类型。很多同学以为数组在 CSV 里也用逗号分隔结果是 Flink 默认的csv.array-element-delimiter是;如果你的数据是a,b,c要么在配置里明确改成csv.array-element-delimiter ,要么在源文件侧统一为分号。还有一点CSV 对 MAP、ROW 这类复杂类型的表达能力很弱一旦 Schema 里出现了复杂嵌套类型解析策略很容易失控我一般会避免。3.2 核心参数速查带 csv. 前缀才生效CSV Format 的参数很多但常用的也就下面这些。注意所有参数在 DDL 的 WITH 里都要带csv.前缀这是新手最容易忽略的点。配置项默认值作用csv.field-delimiter,列分隔符支持csv.quote-character引号字符包裹含特殊字符的字段csv.disable-quote-characterfalse如果置 true要求输入文件不带引号csv.escape-character无转义字符自定义转义规则时用csv.allow-commentsfalse设为 true 后# 开头行视为注释跳过csv.array-element-delimiter;数组元素分隔符csv.null-literal空字符串哪种字面量代表 null比如 NULL、\Ncsv.ignore-parse-errorsfalse解析失败时字段置 null不整体抛错csv.date-formatyyyy-MM-ddDATE 类型解析格式csv.time-formatHH:mm:ssTIME 类型解析格式csv.timestamp-formatyyyy-MM-dd HH:mm:ssTIMESTAMP 类型解析格式csv.timestamp-format.standardSQLSQL 或 ISO-8601决定时间戳格式标准csv.write-bigdecimal-in-scientific-notationfalse写入 BigDecimal 时是否用科学计数法除了这些还需要记住一个原则参数值不是随便抄参数名也不能随手改。比如我曾经见过有人把csv.field-delimiter写成csv.fields-terminator结果参数没生效文件按逗号解析但业务要求竖线分隔整个作业跑出来的结果错得离谱。带csv.前缀是 CSV Format 统一约定不是我的个人偏好。3.3 时间格式是解析重灾区时间格式的坑在 CSV 接入里占比能超过一半。Flink 默认的时间戳格式是 SQL 风格yyyy-MM-dd HH:mm:ss这和你平时在数据库里看到的一致。但现实业务里经常冒出来三种变体第一种是带毫秒的2024-01-01 12:30:00.123。这个默认格式解析不了需要在 WITH 里指定csv.timestamp-format yyyy-MM-dd HH:mm:ss.SSS第二种是斜杠日期2024/01/01 12:30:00。如果字段声明成 TIMESTAMP需要把整个时间格式调成csv.timestamp-format yyyy/MM/dd HH:mm:ss。这里注意不要只改csv.date-format因为csv.date-format只管 DATE 类型的字段TIMESTAMP 字段不受它控制。第三种是 ISO-8601 风格2024-01-01T12:30:00Z。这种带 T 和时区的时间戳SQL 风格的格式解析不了需要把标准切到 ISO-8601声明字段类型为TIMESTAMP_LTZcsv.timestamp-format.standard ISO-8601这里有个容易混淆的点csv.timestamp-format.standard和csv.timestamp-format是两个不同参数。前者决定整体遵循 SQL 还是 ISO-8601 规则后者是具体的 pattern。默认是SQL。当一个 CSV 文件里有多种时间格式混在一起时你没法用一个参数解决所有列需要提前在上游统一格式或者对不同列分别用不同的表定义再 join这属于麻烦但可行的兜底方案。3.4 null 语义与容错策略CSV 文件的空值表达方式千奇百怪有的是空字符串有的是 NULL 四个字母有的是\N有的是 \N。默认情况下CSV Format 会把空字符串识别为 null。如果你希望某个特定字符串也被识别成 null比如数据库导出的文件里用\N表示空值就配置csv.null-literal \N注意这里写的是字面量比较不是正则表达式。每列只要文本内容和配置完全一致都会被当成 null 处理。反过来的情况也存在比如客户端导出 CSV 时把空字符串原样保留而你希望空字符串还是字符串类型那就只能保证字段类型是 STRING并且不要用会触发 null 转换的配置。容错方面最有用的参数是csv.ignore-parse-errors。默认 false 时解析任何一行失败都可能导致作业失败设为 true 时解析失败的那个字段会被置为 null但这一行数据本身仍然会进入下游。注意“字段置 null”和“整行丢弃”不是一回事如果你下游依赖这个字段做非空校验最终还是会报错。我的经验是开发调试阶段保持严格模式false用报错把脏数据暴露出来确认数据质量稳定后再针对个别已知脏字段开启true不要一上来就开容错否则问题会被静默吞掉。4. 实战CSV 文件整库同步到 MySQL / ClickHouse4.1 同步管道怎么设计把 CSV 接入数据库是很多团队的刚需尤其当 CSV 文件来自业务报表、第三方工具或者系统导出时。这里我用一个“CSV 源表同步到 MySQL”的例子说明ClickHouse 的用法几乎一致只要把 JDBC URL、驱动和表结构对应替换就行。源表沿用前面的user_event_csv目标 MySQL 表 DDL 如下CREATE TABLE mysql_sink ( user_id BIGINT, user_name VARCHAR(100), score DECIMAL(10, 2), is_vip TINYINT, event_time TIMESTAMP(3), PRIMARY KEY (user_id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://localhost:3306/test?serverTimezoneAsia/ShanghaiuseSSLfalse, table-name user_event, username root, password yourpassword, sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 2s );注意几个细节serverTimezoneAsia/Shanghai是 MySQL 8 驱动下常见的必填参数不加会直接报时间区错误。useSSLfalse是测试环境常用配置生产环境按安全要求来。sink.buffer-flush.max-rows和sink.buffer-flush.interval控制 JDBC Sink 的批量写入节奏不要设置成 1否则性能会非常差也不要太大否则攒在内存里的行数过多一旦任务失败重放压力会集中到目标库。4.2 调度执行与并行度控制实际跑同步任务时一条 INSERT INTO SELECT 就完成了全流程INSERT INTO mysql_sink SELECT user_id, user_name, TRIM(CAST(score AS STRING)) AS score, CAST(is_vip AS TINYINT), event_time FROM user_event_csv WHERE event_time IS NOT NULL;CLEAN 逻辑可以在 SQL 里直接做。比如WHERE event_time IS NOT NULL过滤掉时间字段缺失的行TRIM去除两侧空格。这些操作在 PyFlink 里全部通过 SQL 完成不需要写 UDF。并行度方面有几个经验。Filesystem 连接器读取 CSV 文件时单个文件的并行度有限如果你处理的是一堆小文件建议把作业并行度设成和文件数量相近的值避免任务倾斜。PyFlink 里通过t_env.get_config().set_parallelism(4)来设置。写入 JDBC 时并行度不要太高因为目标库的连接数和写入压力会成为瓶颈。合理做法是源头并行高sink 端通过sink.buffer-flush参数把写入合并成批这样整体吞吐会比较稳。4.3 JDBC 连接器异常排查实录热词里那波“flink 的 jdbc 连接器异常”其实就是这几个问题的高频集合现象常见原因解决办法ClassNotFoundExceptionJDBC 驱动 jar 没放进来把 mysql-connector-j / clickhouse-jdbc 放进 lib 或依赖Communications link failureurl、端口、host 配置不对先不看 Flink直接用 JDBC 工具测连目标库Unknown time zoneMySQL 连接串缺少 serverTimezone加 serverTimezoneAsia/ShanghaiData truncation目标表字段长度小于源字段检查 VARCHAR 长度、DECIMAL 精度Buffer flush 阶段任务慢目标库连接不够或写入参数不合理调低写入并行度调大 flush 间隔我之前遇到过最典型的一个问题本地 JDBC 工具连接 MySQL 完全正常Flink 作业一跑就报连接超时。排查到最后发现不是代码问题而是 jar 版本冲突导致驱动类被其他版本覆盖。那种情况下你会看到一连串奇怪的 NoSuchMethodError而不是直接的连接失败。解决办法就是统一 Flink、MySQL 驱动和相关依赖的版本不要把不同版本的 connector 混在同一个 classpath 里。4.4 Spring Boot 整合 Flink 的依赖与坑“Spring Boot 整合 Flink”很多人理解成把 Flink 跑在 Spring Boot 进程里这个方向风险很大。我一般推荐的做法是Spring Boot 只负责提交和管理任务Flink 集群独立运行。如果你确实要在工程里嵌入 Flink Table API依赖要配全dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridge/artifactId version你的Flink版本/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-csv/artifactId version你的Flink版本/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc/artifactId version你的Flink版本/version /dependencyJava 代码里用 StreamTableEnvironment 注册和执行 SQL流程和 PyFlink 完全一致。要说坑最大的是版本冲突Spring Boot 内置的依赖管理和 Flink 自带的依赖经常打架尤其是 jackson、guava、netty 这些老熟人。我的经验是 Spring Boot 工程里不要直接依赖 flink-clients 全家桶尽量用 flink-table-api 层运行时把完整 jar 交给 Flink 集群这样冲突面小得多。5. 高频问题与避坑速查5.1 常见问题一张表这里把我在不同项目里真刀真枪撞过的问题汇总成表方便你排查时一条条对照。现象原因解决方法第一行被当成数据CSV Format 不支持表头预处理去掉表头或用通配符绕过数字字段解析失败类型不匹配如千分位、空串检查样例调整字段类型或配置容错中文乱码文件不是 UTF-8 编码先转码或在上游统一编码TIMESTAMP 解析失败时间格式与默认 pattern 不一致配置 csv.timestamp-format日期时间少一天时区处理不一致明确使用 TIMESTAMP 还是 TIMESTAMP_LTZ数组字段读不出默认用分号分隔配置 csv.array-element-delimiternull 变成字符串或反之null-literal 未配置配置 csv.null-literal写入 BigDecimal 变科学计数法写侧参数未设置设置 write-bigdecimal-in-scientific-notationfalseCould not find any factory for identifier csv缺 flink-csv jar添加依赖或放进 lib 目录WITH 参数不生效参数名少了 csv. 前缀逐个检查参数名这张表里的问题有好几个是在生产环境踩过之后才彻底明白的。尤其是“第一行被当成数据”这个光靠代码看不出来因为源文件在本地看明明有表头但 Flink 读取时把它当普通行等你会发现数据量突然多了几行已经晚了。5.2 一份可以直接抄的生产级配置模板下面这份 DDL 基本覆盖了日常 90% 的 CSV 读取场景可以直接替换字段和路径使用CREATE TABLE csv_source ( id BIGINT, business_name STRING, amount DECIMAL(20, 6), source_type STRING, status INT, tags ARRAYSTRING, create_date DATE, create_time TIME, create_ts TIMESTAMP(3) ) WITH ( connector filesystem, path /data/input/csv/*.csv, format csv, csv.field-delimiter ,, csv.quote-character , csv.allow-comments false, csv.ignore-parse-errors false, csv.null-literal , csv.date-format yyyy-MM-dd, csv.time-format HH:mm:ss, csv.timestamp-format yyyy-MM-dd HH:mm:ss, csv.array-element-delimiter ; );如果你拿到的是带引号但没有解析出来的字段先检查quit-character是否被改成其他字符。如果处理的是报表系统导出的文件建议先拿几行样例在 Excel 或者文本编辑器里打开确认分隔符到底是逗号、Tab 还是分号再决定配置项。别小看这一步很多“数据错位”最后都发现是分隔符猜错了。5.3 我在实际工程里的三个小习惯第一每接一个新 CSV先用head -n 5看样例把每一列类型手工标注一遍再对着 DDL 逐列核对。这比在 Flink 里反复跑作业试错快得多。第二开发阶段永远先用严格模式跑小文件让报错把脏数据暴露出来等确认大部分数据干净后再针对个别字段开容错。第三凡是做同步任务sink 表的字段顺序和类型一定要和 SELECT 列表一一对应别依赖“名称相同就自动匹配”CSV 的下游不会帮你纠正列错位。最后再说一个我自己的体会CSV 看起来是所有格式里最没技术含量的一种但它的脏数据问题很考验基本功。遇到问题先别急着写代码把样例数据打开看几行再回到 Schema 上找原因往往比加各种容错参数要省事得多。
返回列表