DTStack / DTStack/chunjun

binlog异步处理机制会导致flink异常中断时数据丢失

Aperta
#533 6 commenti 0 reazioni 0 assegnatari Vedi su GitHub
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

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.