Spanner表/列锁导致Dataflow流式作业严重卡顿问题
resultset.next(),超时远超预期 我们在Dataflow中使用Apache Beam的BatchReadOnlyTransaction.executeQuery(Statement)执行只读事务读取数据时遇到异常:当目标数据的部分列存在排他锁时,锁等待时间仅约2秒,但Dataflow作业步骤却无输出运行了20分钟,日志显示卡在resultset.next(),实际请求耗时近20分钟才完成,完全不符合“锁释放或等待2秒后继续推进”的预期。
代码示例
try (ResultSet rs = readOnlyTransaction.executeQuery(statement)) { while (rs.next()) { Struct struct = rs.getCurrentRowAsStruct(); } }
捕获到的异常日志(已翻译)
步骤**中的操作已持续至少30分钟未输出或完成,当前处于处理状态
at java.base@11.0.9/jdk.internal.misc.Unsafe.park(Native Method)
at java.base@11.0.9/java.util.concurrent.locks.LockSupport.park(LockSupport.java:194)
at java.base@11.0.9/java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await(AbstractQueuedSynchronizer.java:2081)
at java.base@11.0.9/java.util.concurrent.LinkedBlockingQueue.take(LinkedBlockingQueue.java:433)
at app//com.google.cloud.spanner.AbstractResultSet$GrpcStreamIterator.computeNext(AbstractResultSet.java:939)
at app//com.google.cloud.spanner.AbstractResultSet$GrpcStreamIterator.computeNext(AbstractResultSet.java:884)
at app//com.google.common.collect.AbstractIterator.tryToComputeNext(AbstractIterator.java:146)
at app//com.google.common.collect.AbstractIterator.hasNext(AbstractIterator.java:141)
at app//com.google.cloud.spanner.AbstractResultSet$ResumableStreamIterator.computeNext(AbstractResultSet.java:1135)
at app//com.google.cloud.spanner.AbstractResultSet$ResumableStreamIterator.computeNext(AbstractResultSet.java:1007)
at app//com.google.common.collect.AbstractIterator.tryToComputeNext(AbstractIterator.java:146)
at app//com.google.common.collect.AbstractIterator.hasNext(AbstractIterator.java:141)
at app//com.google.cloud.spanner.AbstractResultSet$GrpcValueIterator.ensureReady(AbstractResultSet.java:270)
at app//com.google.cloud.spanner.AbstractResultSet$GrpcValueIterator.getMetadata(AbstractResultSet.java:246)
at app//com.google.cloud.spanner.AbstractResultSet$GrpcResultSet.next(AbstractResultSet.java:120)
问题分析
从日志栈可以看出,代码卡在LinkedBlockingQueue.take()方法,说明Spanner的GRPC流迭代器在等待数据时进入了无限阻塞状态,没有触发预设的超时机制。核心原因可能是:
- 只读事务未配置超时参数,导致遇到锁等待时不会主动中断请求
- GRPC流的读取超时未生效,迭代器持续等待数据而不抛出异常
- Dataflow作业未设置步骤级超时,卡住的步骤无法自动终止重试
解决方案
- 为只读事务设置全局超时
在创建BatchReadOnlyTransaction时,通过TransactionOptions明确设置超时时间(建议覆盖锁等待时间,留足缓冲):
BatchReadOnlyTransaction transaction = spannerClient.batchReadOnlyTransaction( TransactionOptions.newBuilder() .setReadOnly(ReadOnly.newBuilder().setTimestampBound(TimestampBound.strong())) .setTimeout(Duration.ofSeconds(20)) .build() );
- 为查询设置单独超时
执行查询时,通过Statement的setQueryTimeout方法设置查询级别的超时:
Statement statement = Statement.newBuilder(sql) .setQueryTimeout(Duration.ofSeconds(15)) .build();
- 配置Dataflow作业的超时与重试策略
在Pipeline配置中设置作业超时和重试策略,确保卡住的步骤能被自动终止并重试:
PipelineOptions options = PipelineOptionsFactory.create(); DataflowPipelineOptions dataflowOptions = options.as(DataflowPipelineOptions.class); dataflowOptions.setJobTimeout(Duration.ofMinutes(30)); dataflowOptions.setRetryStrategy(RetryStrategy.RETRY_ON_TRANSIENT_ERROR);
- 排查Spanner锁的实际持有情况
通过Spanner控制台查看排他锁的持有事务时长,确认是否存在长事务未释放锁的情况,这可能是导致只读事务持续等待的根本原因。
内容的提问来源于stack exchange,提问作者Keshav

