ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink SQL MySQl数据同步到ES数量未对齐的问题

Flink SQL MySQl数据同步到ES数量未对齐的问题 目录问题背景排除过程排除出未同步的商品数据打开ES新增修改日志打开Flink debug日志发现解决方案参考问题背景使用Flink SQL把MySQL的商品数据同步到ES发现ES的数据条数总是会少一些。SQL大概是这样的-- 商品数据以及这个商品对应实物最小克重INSERTINTOes_productSELECTmp.id,mp.name,agg.minWeight,...省略...FROMmysql_productASmpLEFTJOIN(SELECTproduct_id,MIN(weight)ASminWeight...省略...FROMmysql_product_codeWHEREcode_status0GROUPBYproduct_id)ASaggONagg.product_idmp.id...省略...排除过程排除出未同步的商品数据我写了一个接口查询全部表数据一条条向ES查询判断有没有找出未同步的数据。打开ES新增修改日志logger.index_indexing_slowlog.name org.elasticsearch.index.indexing.slowlog logger.index_indexing_slowlog.level TRACE logger.index_indexing_slowlog.appenderRef.index_indexing_slowlog.ref index_indexing_slowlogPUT /product_dev/_settings { index.search.slowlog.threshold.query.warn: 10s, index.search.slowlog.threshold.fetch.warn: 10s, index.indexing.slowlog.threshold.index.warn: 0ms, index.indexing.slowlog.threshold.index.info: 0ms, index.indexing.slowlog.threshold.index.debug: 0ms, index.indexing.slowlog.threshold.index.trace: 0ms, index.indexing.slowlog.source: true }排查完后关闭慢日志PUT /product_dev/_settings { index.indexing.slowlog.threshold.index.warn: -1, index.indexing.slowlog.threshold.index.info: -1, index.indexing.slowlog.threshold.index.debug: -1, index.indexing.slowlog.threshold.index.trace: -1, index.indexing.slowlog.source: false }为了防止日志文件被压缩或清理config/log4j2.propertiesappender.rolling.policies.size.size 111128MB打开Flink debug日志设置环境变量 ,编辑~/.bashrc,加入一下内容exportROOT_LOG_LEVELDEBUGexportMAX_LOG_FILE_NUMBER1000source~/.bashrc为了防止日志文件被压缩或清理conf/log4j.propertiesappender.main.policies.size.size 10000MB发现针对未同步到ES的数据查看Flink日志有向ES发送并且ES日志也有打印这些数据。豆包与Gemini都提到我的SQL双流JOIN 回撤的问题。修改我的日志打印。这块代码是我之前写的 flink-sql-connector-elasticsearch8Overridepublicvoidinvoke(RowDatavalue,Contextcontext)throwsException{// 获取当前行操作类型RowKindkindvalue.getRowKind();StringdocIdfieldGetters[primaryKeyIndex].getFieldOrNull(value).toString();log.info(write ES, index{}, docId{}, op{},index,docId,kind);if(kindRowKind.INSERT||kindRowKind.UPDATE_AFTER){// Upsert 操作加入批量缓冲不阻塞当前线程吞吐远高于逐条写MapString,ObjectdocrowToMap(value);bulkIngester.add(op-op.index(i-i.index(index).id(docId).document(doc)));log.debug(write ES success, index{}, docId{}, doc{},index,docId,doc);}elseif(kindRowKind.DELETE){// 删除操作同样进入批量与写入共用同一缓冲队列bulkIngester.add(op-op.delete(d-d.index(index).id(docId)));}// UPDATE_BEFORE 通常忽略因为紧接着的 UPDATE_AFTER 会覆盖整个 Doc}果然发现有DELETE操作。解决方案攒批、去重---------------------------------------------------------------------------------- 1. MiniBatch 微批攒批解决 ES 频发删改乱序、提升 Sink 吞吐--------------------------------------------------------------------------------SETtable.exec.mini-batch.enabledtrue;SETtable.exec.mini-batch.allow-latency2s;SETtable.exec.mini-batch.size5000;---------------------------------------------------------------------------------- 2. 撤回流与聚合优化消除中间冗余 -D/-U 消息、防止 GROUP BY 热点---------------------------------------------------------------------------------- 在微批内存中抵消/去重同一批次内的 -D / -U / U只输出最终状态SETtable.exec.retract.deduplicatetrue;-- 开启 Local-Global 两阶段聚合解决热点 Key 倾斜降低 RocksDB I/OSETtable.optimizer.agg-phase-strategyTWO_PHASE;---------------------------------------------------------------------------------- 3. 状态与物理执行优化---------------------------------------------------------------------------------- 状态过期时间按需调整如 7 天若不设置GROUP BY / JOIN 的状态将永久保存直到内存/磁盘溢出-- SET table.exec.state.ttl 7d;-- 开启子计划复用降低拓扑复杂度SETtable.optimizer.reuse-sub-plan-enabledtrue;兜底我的业务没有物理删除商品数据索引屏蔽DELETE操作。Overridepublicvoidinvoke(RowDatavalue,Contextcontext)throwsException{// 获取当前行操作类型RowKindkindvalue.getRowKind();StringdocIdfieldGetters[primaryKeyIndex].getFieldOrNull(value).toString();log.info(write ES, index{}, docId{}, op{},index,docId,kind);if(kindRowKind.INSERT||kindRowKind.UPDATE_AFTER){// Upsert 操作加入批量缓冲不阻塞当前线程吞吐远高于逐条写MapString,ObjectdocrowToMap(value);bulkIngester.add(op-op.index(i-i.index(index).id(docId).document(doc)));log.debug(write ES success, index{}, docId{}, doc{},index,docId,doc);}elseif(kindRowKind.DELETE){// 删除操作同样进入批量与写入共用同一缓冲队列// 双流join撤回流问题一般不会物理删除先屏蔽// bulkIngester.add(op - op.delete(d - d.index(index).id(docId)));}// UPDATE_BEFORE 通常忽略因为紧接着的 UPDATE_AFTER 会覆盖整个 Doc}参考豆包、Gemini
返回列表