Flink Print SQL Connector最强“肉眼调试”Sink,用对真的省一半时间
·
1. 最小可用 DDL
CREATE TABLE print_table (
f0 INT,
f1 INT,
f2 STRING,
f3 DOUBLE
) WITH (
'connector' = 'print'
);
然后:
INSERT INTO print_table
SELECT f0, f1, f2, f3 FROM some_table;
2. 输出长什么样(RowKind 很关键)
输出格式固定是:
$row_kind(f0,f1,f2…)- 例如:
+I(1,1)表示 插入 一条记录
常见 RowKind(你排障时最该关注):
+I:INSERT(新增)-D:DELETE(删除)-U:UPDATE_BEFORE(更新前)+U:UPDATE_AFTER(更新后)
如果你做的是聚合、TopN、Join 之类的 changelog 查询,看到 -U/+U 是正常的;如果你期望 append-only,却看到 update/delete,就说明上游不是你以为的模式。
3. print-identifier:多任务并行时救命
并行度 > 1 时,不同 subtask 会同时打印,日志会混在一起。print-identifier 用来给每行前面加一个前缀,方便 grep。
CREATE TABLE print_table (
f0 INT,
f1 STRING
) WITH (
'connector' = 'print',
'print-identifier' = 'DEBUG-PIPE'
);
并行度和 identifier 组合后的展示规律
- parallelism > 1:会有 taskId 信息(用于区分 subtask)
- parallelism == 1:只有一条打印流,taskId 可省略
- 提供 print-identifier:会把 identifier 作为前缀
- 不提供:就只打印 taskId 或纯 output
(你贴的那张表就是在解释这个组合效果。)
4. standard-error:把输出打到 stderr(生产更好用)
很多日志采集/告警系统对 stderr 会单独处理,排障时更容易筛出来:
CREATE TABLE print_err (
f0 INT,
f1 STRING
) WITH (
'connector' = 'print',
'standard-error' = 'true',
'print-identifier' = 'ERR-TRACE'
);
5. sink.parallelism:控制打印并行度(很实用)
默认并行度跟随上游,有时上游并行度很高会刷屏。你可以主动压低:
CREATE TABLE print_table (
f0 INT,
f1 STRING
) WITH (
'connector' = 'print',
'sink.parallelism' = '1',
'print-identifier' = 'ONE'
);
调试时经常这么干:打印单并行,数据更好读。
6. 最常见的组合:DataGen + Print(1 分钟验证 SQL)
你上一条刚讲了 DataGen,这里给一套“直接跑”的模板:
CREATE TABLE gen_src (
id BIGINT,
score INT,
name STRING
) WITH (
'connector' = 'datagen',
'rows-per-second' = '5',
'fields.id.kind' = 'sequence',
'fields.id.start' = '1',
'fields.id.end' = '20',
'fields.score.min' = '0',
'fields.score.max' = '100',
'fields.name.length' = '10',
'fields.name.var-len' = 'true'
);
CREATE TABLE print_sink (
id BIGINT,
score INT,
name STRING
) WITH (
'connector' = 'print',
'print-identifier' = 'GEN-OUT',
'sink.parallelism' = '1'
);
INSERT INTO print_sink
SELECT * FROM gen_src;
你会在 task 日志里看到类似:
+I(1,87,abc...)+I(2,12,xyz...)
7. 生产排障小技巧(不踩坑版)
- 别长期开:Print 会大量 IO,可能把 TaskManager 日志打爆
- 加过滤:先 WHERE 限定某个 user_id / trace_id 再打印
- 单并行:
sink.parallelism=1让日志可读 - 打印关键字段:用
SELECT id, op, ts, ...,别把大字段全打出来 - 注意看 RowKind:排查“为什么结果会跳、会回撤”时,RowKind 是第一线索
更多推荐
所有评论(0)