paimon表producer mode必须配置对.
一般建议lookup,上游是binlog则配置input , paimon ods层或者append table表配置none即可 详细选择看另一篇帖子
如果1 上游表是主键表,2 表无法提供-U 即你配置producer mode=none 导致没有-U,3 下游需要retract语义(比如下游sum聚合,比如统计订单金额,有个金额修正,你不给-U +U,只给+U 那么统计结果就不对; 2 join场景 3 下游sink表需要retract场景 ). 那么下游消费该paimon的flink会自动补充ChangelogNormalize算子,如果多个flink消费,那么重复的lookup压力给到了flink,不如直接配置producer mode = lookup
流式概念 | Apache Flink
paimon聚合表做状态,而不是flink做状态
// flink升级,业务迭代等可能导致状态丢失,此场景用paimon聚合表更合理方便,sr表也可以,sr稳定性略差
paimon配置 lookup-wait=true //默认就是true
可能为了异步compact,会配置为false,这会导致flink failover 时 changelog 还没生成完
paimon 表使用agg 注意配置table.exec.sink.upsert-materialize=NONE
paimon文档有写
即: flink sink paimon agg表,如果你不配置为NONE,flink大概率会误认为该接受retract流的piamon需要upsert流, 举例数据案例:
这会导致最终值错误,是30而不是20,导致agg表中last_value统计错误;
也可能导致值重发,比如+30发了两次,会导致sum这种聚合函数重复统计,数值变大
paimon表使用了 agg类函数 注意 sequence.field 的设置
如果配置的字段 sequence.field 可能存在乱序,会导致认为是迟到数据,不去更新agg函数的值,是否影响业务正确性
paimon agg function并非都支持 retraction回撤
Aggregation | Apache Paimon
有的函数支持
有的建议关闭回撤消息感知
有的不兼容回撤消息,会导致函数结果错误,这个参考最新文档处理
为了避免并发写入冲突 建议metastore用jdbc且开启并发 lock.enable
// flink实时写入 + 离线批修正数据
paimon实时数仓 如何防止表提前过期快照导致下游消费异常
// 下游消费延迟,对应快照被上游paimon给快照过期了,丢数据
相关参数:
paimon ddl:
'snapshot.time-retained' = '7d', -- 快照超过7天就过期
'snapshot.num-retained.min' = '20',
'snapshot.num-retained.max' = '50' -- 最多保留最近的50个快照,flink写入chk太频繁很容易达到
'consumer.expiration-time' = '7d' -- 如果 consumer 长期不更新,超过 consumer.expiration-time,Paimon 会认为它是死 consumer,清理 consumer 记录,不再保护旧 snapshot
注意:
snapshot.num-retained.min 是硬下限
snapshot.num-retained.max 是硬上限
snapshot.time-retained 是时间窗口条件
flink read paimon参数:
'consumer-id' = 'xxx',
'consumer.mode' = 'exactly-once'
paimon小文件问题
写入前规避、写入中调参、写入后治理
写入前规避
分区 分桶不要太细
写入中调参
flink chk不要太频繁
flink sink并行度不要太高
控制 changelog-producer lookup 会额外生成 changelog 文件。只有下游确实需要完整 -U/+U 时才开
关闭不必要的 upsert materialize
开启 precommit compact:
如果 changelog 小文件很多,可以考虑:'precommit-compact' = 'true' 作用是在提交前合并小 changelog 文件 可能导致chk变慢
优化compact参数:
'compaction.min.file-num' = '5',
'compaction.max.file-num' = '50',
'num-sorted-run.compaction-trigger' = '10'
写入后治理
离线 compact 历史分区