binlog异步处理机制会导致flink异常中断时数据丢失
Open
- Dominant language
- Java
- Stars
- 4.1k
- Forks
- 1.7k
- PR merge metrics
- No merged PRs in 30d
Description
版本:1.12-release
canal的EventParser在每个transaction的`sink`后会获取当前事物的position,然后persistLogPosition持久化点位信息。
flinkx实现了canal的logPositionManager,persistLogPosition会做2个动作:
1. flink状态更新:format.setEntryPosition(logPosition.getPostion())
2. 本地缓存更新:logPositionCache.put(destination, logPosition);
flinkx的binlogReader的实现中,canal的sink是通过queue异步处理的,由flink的`DtInputFormatSourceFunction`在执行`nextRecord`时从queue中poll出来处理row。
那么问题来了,
如果此时server断电flink程序异常中断,format的state已经往前走,但是异步处理比较慢,还没处理完被异常中断了,
重启时读取的checkpoint的点位是后面的position,会导致有些日志数据未被处理。
Contributor guide
Assessment
This issue has not been assessed yet.