binlog异步处理机制会导致flink异常中断时数据丢失
- Lingua principale
- Java
- Stelle
- 4.1k
- Fork
- 1.7k
- Metriche di merge delle PR
- Nessuna PR unita negli ultimi 30g
Descrizione
版本: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,会导致有些日志数据未被处理。
Guida per i contributori
Apri la guida per i contributori
Direzione di ricerca
Traccia il flusso di EventParser e logPositionManager di canal fino a binlogReader, in particolare persistLogPosition, format.setEntryPosition, logPositionCache e DtInputFormatSourceFunction.nextRecord. Riproduci un'interruzione mentre rimangono record accodati non elaborati, quindi definisci il completamento come un recupero dopo il riavvio che non salti tali record e copri il comportamento con un test di regressione.
Scritto dal modello di indicizzazione a partire dal testo della issue.
Valutazione
- Stack tecnologico
- java
- Ambito
- data-engineering, stream-processing
- Tipo di issue
- Bug
- Difficoltà
- 4/5
- Tempo stimato
- 3-5 giorni
- Stato di attività
- Ferma
- Chiarezza
- Da chiarire
- Idoneità per principianti
- 30/100