wangbing505 opened a new issue, #11889:
URL: https://github.com/apache/seatunnel/issues/11889

   ### Search before asking
   
   - [x] I had searched in the 
[issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue+label%3A%22bug%22)
 and found no similar issues.
   
   
   ### What happened
   
   使用非精准一致性抽取,快照阶段数据抽取到 120,430 条后停止增长,作业完全卡死,还处于快照阶段
   
   ### SeaTunnel Version
   
   使用的版本是2.3.13 zeta引擎
   
   ### SeaTunnel Config
   
   ```conf
   source {
     MongoDB-CDC {
       hosts = "..."
       username = "..."
       password = "***"
       database = ["ai-correction-solution"]
       collection = [
         "ai-correction-solution.originalResponse",
         "ai-correction-solution.keypointResponse",
         "ai-correction-solution.originalRequest",
         "ai-correction-solution.subTaskRequest",
         "ai-correction-solution.subTaskResponse"
       ]
       batch.size = 8192
       poll.max.batch.size = 8192
       incremental.snapshot.chunk.size.mb = 64
       startup.mode = "initial"
       exactly_once = false  # 关键配置
       tables_configs = [
         {
           schema {
             table = "***.***"
             fields {
               "topic"  : string,
               "key"    : string,
               "value"  : string
             }
           }
         }]
       # ... 其他配置
     }
   }
   
   
   
   sink {
    console {
     }
   }
   ```
   
   ### Running Command
   
   ```shell
   local模式运行
   ```
   
   ### Error Exception
   
   ```log
   现象:快照阶段数据抽取到 120,430 条后停止增长,作业完全卡死,还处于快照阶段
   ```
   
   ### Zeta or Flink or Spark Version
   
   作业最后一次有进展的日志
   `2026-08-19 10:02:35,655 INFO - Snapshot step 1 - Determining low watermark 
     {resumeToken={"_data": "826A850EB3000003ED2B0429296E1404"}, 
timestamp=7675557301884814317}
   
   2026-08-19 10:02:35,655 INFO - Snapshot step 3 - Determining high watermark 
     {resumeToken={"_data": "826A850EBA000000E72B0429296E1404"}, 
timestamp=7675557331949584615}`
   
   二、问题排查过程
   2.1 Phase 1:线程栈分析
   对运行中的作业进行了两次 jstack 线程转储(间隔 5 分钟),两次结果完全一致,关键线程状态如下:
   关键线程 #306:debezium-snapshot-reader-0
   状态: WAITING (parking)
   `java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
   
java.util.concurrent.locks.AbstractQueuedSynchronizer.parkAndCheckInterrupt(AbstractQueuedSynchronizer.java:836)
   
java.util.concurrent.locks.AbstractQueuedSynchronizer.doAcquireSharedInterruptibly(AbstractQueuedSynchronizer.java:997)
   
java.util.concurrent.locks.AbstractQueuedSynchronizer.acquireSharedInterruptibly(AbstractQueuedSynchronizer.java:1304)
   java.util.concurrent.Semaphore.acquire(Semaphore.java:312)
   
java.util.concurrent.LinkedBlockingDeque.putLast(LinkedBlockingDeque.java:396)
   
io.debezium.connector.base.ChangeEventQueue.doEnqueue(ChangeEventQueue.java:237)
   
org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.source.fetch.MongodbStreamFetchTask.execute(MongodbStreamFetchTask.java:203)
   
org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.source.fetch.MongodbScanFetchTask.execute(MongodbScanFetchTask.java:151)`
   诊断: 快照读取线程阻塞在向 ChangeEventQueue(队列①)写入数据,队列已满无法继续。
   关键线程 #284:Source Data Fetcher
   状态: TIMED_WAITING (parking)
   `java.util.concurrent.locks.LockSupport.parkNanos(LockSupport.java:215)
   
java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.awaitNanos(AbstractQueuedSynchronizer.java:2078)
   java.util.concurrent.LinkedBlockingQueue.offer(LinkedBlockingQueue.java:385)
   
org.apache.seatunnel.connectors.seatunnel.common.source.reader.fetcher.FetchTask.run(FetchTask.java:59)`
   诊断: 数据获取线程阻塞在向 elementsQueue(队列②)写入数据,队列也已满。
   关键线程 #261:BlockingWorker taskGroupId=3
   状态: TIMED_WAITING (sleeping)
   `java.lang.Thread.sleep(Native Method)
   
org.apache.seatunnel.connectors.seatunnel.common.source.reader.SourceReaderBase.getNextFetch(SourceReaderBase.java:172)
   
org.apache.seatunnel.connectors.seatunnel.common.source.reader.SourceReaderBase.pollNext(SourceReaderBase.java:150)`
   诊断: Source 任务线程从队列② poll 数据返回 null,进入 100ms sleep 空转循环。
   下游线程状态
   - HDFS DataStreamer: WAITING 在 Object.wait,无数据可写
   初步结论
   1. 无 JVM 级死锁: jstack 未报告 "Found one Java-level deadlock"
   2. 上游队列①②都已满: 生产者阻塞在 put/offer 操作
   3. 消费者认为无数据: 线程 #261 的 poll 返回 null,进入 sleep 循环
   4. 下游全部饿死: Sink 线程都在等待数据
   核心矛盾: 队列②既然已满(线程 #284 阻塞在 offer),为什么线程 #261 的 poll 会返回 null?
   2.2 Phase 2:消费者源码分析
   源码:SourceReaderBase.getNextFetch() (行 166-181)
   `RecordsWithSplitIds<E> recordsWithSplitId = elementsQueue.poll();
   if (recordsWithSplitId == null || !moveToNextSplit(recordsWithSplitId, 
output)) {
       Thread.sleep(100);
       return null;
   }`
   发现: poll 是非阻塞的,sleep 的原因有两种:
   1. poll 返回 null(队列为空)
   2. moveToNextSplit 返回 false(没有有效 split 可推进)
   继续追踪 elementsQueue 的数据来源 → IncrementalSourceScanFetcher.pollSplitRecords()
   源码:IncrementalSourceScanFetcher.pollSplitRecords() (行 119-127)
   `@Override
   public RecordsWithSplitIds<SourceRecord> pollSplitRecords() throws 
InterruptedException {
       checkReadException();
       if (hasNextElement.get()) {
           return pollSplitRecordsIfNotExactlyOnce();
       }
       return null;  // ← hasNextElement=false 后直接返回 null
   }`
   源码:IncrementalSourceScanFetcher.pollSplitRecordsIfNotExactlyOnce() (行 
130-146)
   `private RecordsWithSplitIds<SourceRecord> 
pollSplitRecordsIfNotExactlyOnce() {
       List<SourceRecord> sendRecords = new ArrayList<>();
       List<DataChangeEvent> batch = queue.poll();  // ← 从队列① 取数据
       
       if (batch != null) {
           for (DataChangeEvent event : batch) {
               SourceRecord record = event.getRecord();
               sendRecords.add(record);
               
               // ← 关键:读到 HIGH watermark 就停止消费
               if (isHighWatermarkEvent(record)) {
                   hasNextElement.set(false);
               }
           }
       }
       
       return new RecordsWithSplitIds<>(sendRecords);
   }`
   关键发现: 消费者在读到 HIGH watermark 事件后,立即设置 hasNextElement=false,从此不再从队列① poll 数据。
   矛盾解开
   队列① 中事件的顺序是:
   [LOW watermark] [snapshot 数据 × 120,430] [HIGH watermark] [backfill 变更...] 
[END watermark]
                                                 ↑                    ↑
                                           消费者读到这里就         这些数据无人消费
                                           停止 poll 队列①
   结论: HIGH watermark 之后的 backfill 事件永远不会被消费,导致队列① 被填满,生产者永久阻塞。
   2.3 Phase 3:生产者源码分析
   源码:MongodbScanFetchTask.execute() (行 82-162)
   快照块的执行流程:
   `// Step 1: 发送 LOW watermark (行 94-101)
   changeEventQueue.enqueue(
       new DataChangeEvent(
           WatermarkEvent.create(..., WatermarkKind.LOW, lowWatermark)));
   
   // Step 2: 读取快照数据 (行 107-116)
   MongoCollection<BsonDocument> collection = ...;
   MongoCursor<BsonDocument> cursor = ...;
   while (cursor.hasNext()) {
       BsonDocument document = cursor.next();
       changeEventQueue.enqueue(
           new DataChangeEvent(createSourceRecord(..., document)));
       // ← 120,430 次循环
   }
   
   // Step 3: 发送 HIGH watermark (行 123-130)
   changeEventQueue.enqueue(
       new DataChangeEvent(
           WatermarkEvent.create(..., WatermarkKind.HIGH, highWatermark)));
   
   // Step 4: 执行 backfill (行 135-152) ← 关键!
   final IncrementalSplit dataBackfillSplit = 
       createBackfillStreamSplit(lowWatermark, highWatermark);
   
   final boolean streamBackfillRequired = 
       
dataBackfillSplit.getStopOffset().isAfter(dataBackfillSplit.getStartupOffset());
   
   if (!streamBackfillRequired) {
       changeEventQueue.enqueue(...END watermark...);
   } else {
       // ← 同步执行 backfill,复用同一个队列①
       FetchTask<SourceSplitBase> dataBackfillTask = 
           dialect.createFetchTask(dataBackfillSplit);
       dataBackfillTask.execute(taskContext);  // 行 151
   }`
   关键确认:
   1. 同步执行: backfill 在同一个快照线程内同步执行(第 151 行)
   2. 共享队列: backfill 复用同一个 ChangeEventQueue(队列①)
   3. 顺序保证: HIGH watermark 在 backfill 事件之前就已入队
   非快照执行队列大小由batch.size参数控制 
   
   源码:MongodbStreamFetchTask.execute() bounded read 分支 (行 197-224)
   `if (isBoundedRead()) {  // ← backfill 是 bounded read
       ChangeStreamOffset currentOffset;
       
       if (changeRecord != null) {
           currentOffset = new ChangeStreamOffset(getResumeToken(changeRecord));
           
           // ← 只处理 [startOffset, stopOffset] 窗口内的事件
           if (currentOffset.isAtOrBefore(streamSplit.getStopOffset())) {
               queue.enqueue(new DataChangeEvent(changeRecord));  // 行 203 ← 阻塞点
           }
       } else {
           currentOffset = new 
ChangeStreamOffset(getCurrentClusterTime(mongoClient));
       }
       
       // ← 退出判断
       if (currentOffset.isAtOrAfter(streamSplit.getStopOffset())) {  // 行 211
           queue.enqueue(...END watermark...);
           break;  // 行 222
       }
   }`
   关键问题: 退出判断(行 211)在 enqueue(行 203)之后、同一循环内。
   阻塞场景
   当 backfill 要重放的 [low, high] 窗口内,5 个集合的变更总量超过队列剩余容量(≤ 8192)时:
   1. 线程 #306 在行 203 的 queue.enqueue → LinkedBlockingDeque.putLast 上阻塞(队列满)
   2. 永远走不到行 211 的退出判断
   3. execute() 方法永不返回
   4. snapshotSplitReadTask.isRunning() 恒为 true
   5. IncrementalSourceScanFetcher.isFinished() 恒为 false(判定逻辑:!taskRunning || 
!hasNextElement && reachEnd)
   6. Split 卡死,下一个块无法开始,整个流水线停摆
   
   所以这个bacfill在非精准一致的场景下 需要吗?
   
   ### Java or Scala Version
   
   java8
   
   ### Screenshots
   
   arthas查看 hasNextElement
   
   <img width="996" height="590" alt="Image" 
src="https://github.com/user-attachments/assets/2c95fe50-54c4-4d1a-aca1-19abe43aaf4a";
 />
   
   ### Are you willing to submit PR?
   
   - [ ] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://www.apache.org/foundation/policies/conduct)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to