diff --git a/connectors/starrocks-connector/src/main/java/io/tapdata/connector/starrocks/streamload/StarrocksStreamLoader.java b/connectors/starrocks-connector/src/main/java/io/tapdata/connector/starrocks/streamload/StarrocksStreamLoader.java index 2e60d7d22..85962ea79 100644 --- a/connectors/starrocks-connector/src/main/java/io/tapdata/connector/starrocks/streamload/StarrocksStreamLoader.java +++ b/connectors/starrocks-connector/src/main/java/io/tapdata/connector/starrocks/streamload/StarrocksStreamLoader.java @@ -279,6 +279,9 @@ private void initializeFlushScheduler() { */ private void checkAndFlushIfNeeded() throws StarrocksRetryableException { synchronized (writeLock) { + // 驱逐队首连续的、已全部落库(不在 pendingFlushTables)的陈旧 offset 锚点,避免长期无变更的冷表占据队首,导致断点被永久钉死在其旧位点上 + evictStaleFirstOffsets(); + if (pendingFlushTables.isEmpty()) { return; // 没有数据需要刷新 } @@ -317,6 +320,35 @@ private void checkAndFlushIfNeeded() throws StarrocksRetryableException { } } + /** + * 驱逐 firstOffsetByTable 队首连续的、已全部落库的陈旧 offset 锚点。 + *

+ * flushTable 回传断点时取的是队首表的 offset。若一张长期无变更的冷表在全量/增量早期 + * 进入队首后再无新数据,它就不会再进入 pendingFlushTables,也就永远不会被 flush, + * 从而永远不会被移出队首(见 flushTable 中 tableName.equals(firstTableName) 分支), + * 导致所有其它表 flush 回传的都是这张冷表的旧位点,断点被永久钉死。 + *

+ * 判据:队首表若不在 pendingFlushTables,说明它此刻没有待落库数据、之前的数据均已成功 + * 导入,丢弃其 offset 锚点不会丢数,可直接驱逐;一旦遇到队首表在 pendingFlushTables 中, + * 则停止驱逐——该表会被定时/大小阈值 flush 后按正常流程移出队首并推进断点。 + */ + private void evictStaleFirstOffsets() { + synchronized (firstOffsetByTable) { + Iterator> iterator = firstOffsetByTable.entrySet().iterator(); + while (iterator.hasNext()) { + Map.Entry firstEntry = iterator.next(); + String firstTableName = firstEntry.getKey(); + if (pendingFlushTables.contains(firstTableName)) { + // 队首表仍有待落库数据,等它被正常 flush 后自然推进,停止驱逐 + break; + } + iterator.remove(); + taplogger.info("Evicted stale offset anchor of table {} from queue head " + + "(already flushed, no pending data), allowing breakpoint to advance", firstTableName); + } + } + } + /** * 刷新指定的表 */