从GCP Spanner读取变更流后用JdbcIO查询时遇自动提交异常
Apache Beam + Spanner变更流读取时的Autocommit异常问题
问题现象
使用Apache Beam构建数据管道,从GCP Cloud Spanner读取变更流数据,后续通过JdbcIO.readAll()关联查询另一Spanner表时,间歇性抛出以下异常:
INFO: Autocommit has been disabled[WARNING] org.apache.beam.sdk.Pipeline$PipelineExecutionException: com.google.cloud.spanner.jdbc.JdbcSqlExceptionFactory$JdbcSqlExceptionImpl: FAILED_PRECONDITION: Cannot set autocommit while a transaction is active
同时Spanner连接会自动关闭,已知Spanner JDBC默认autocommit为true,异常原因不明。
核心代码片段
变更流读取主流程
Pipeline p = Pipeline.create(options); PCollection<DataChangeRecord> inputChangeStreams = p.apply(SpannerIO .readChangeStream() .withSpannerConfig(spannerConfig) .withChangeStreamName(options.getChangeStreamTable()) .withMetadataInstance(options.getInputInstanceId()) .withMetadataDatabase(options.getInputDatabaseId()) .withInclusiveStartAt(Timestamp.now())); PCollection<KV<String, String>> enrichedDataStream = buildData(options, inputChangeStreams);
数据关联查询逻辑(buildData方法片段)
private static PCollection<KV<String, String>> buildData(SpannerToSpannerPipelineOption options, PCollection<DataChangeRecord> inputChangeStreams) { PCollection<KV<String, String>> eData = null; PCollection<KV<String, Long>> changeStreamedData = inputChangeStreams .apply("Apply to get Id from CS", ParDo.of(new DoFn<DataChangeRecord, Long>() { @ProcessElement public void processElement(ProcessContext c) { c.element().getMods().forEach(mod -> { JSONObject keysJson = new JSONObject(mod.getKeysJson()); LOG.info("cs keyJson is:" + keysJson); long Id = keysJson.getLong(CONSTANTS.ID); c.output(Id); } ); } })) .apply("Apply Random Key", ParDo.of(new DoFn<Long, KV<String, Long>>() { @ProcessElement public void processElement(ProcessContext c) { c.output(KV.of(UUID.randomUUID().toString(), c.element())); } })).apply("Apply Fixed Length Windows", Window.into(FixedWindows.of(Duration.standardSeconds(1)))); changeStreamedData .apply("Print log", MapElements.into(TypeDescriptor.of(String.class)).via(x -> { LOG.info("incoming data is: " + x.getKey() + "value is" + x.getValue()); return x.getKey(); } )); eData = changeStreamedData .apply("Take each Id and lookup in the Input table", JdbcIO.<KV<String, Long>, KV<String, String>>readAll() .withDataSourceConfiguration( JdbcIO.DataSourceConfiguration.create( CONSTANTS.SPANNERDRIVER, "jdbc:cloudspanner:/projects/" + options.getInputProjectId() + "/instances/" + options.getInputInstanceId() + "/databases/" + options.getInputDatabaseId() ) ) .withQuery(CONSTANTS.QUERY) .withParameterSetter((JdbcIO.PreparedStatementSetter<KV<String, Long>>) (element, preparedStatement) -> { preparedStatement.setLong(1, element.getValue()); LOG.info("element value is: " + element.getValue()); }) .withRowMapper((JdbcIO.RowMapper<KV<String, String>>) resultSet -> { String key = resultSet.getString(CONSTANTS.ID); LOG.info("key value is: " + key); ResultSetMetaData metadata = resultSet.getMetaData(); // Iterate through each column to output Map<String, Object> row = IntStream // Iterate through each column to output .range(1, resultSet.getMetaData().getColumnCount() + 1) // Output a pair with the column name and value. .mapToObj(i -> { try { LOG.info("Metadata2:" + KV.of(metadata.getColumnName(i), (resultSet != null && resultSet.getObject(i) != null) ? resultSet.getObject(i) : "")); return KV.of(metadata.getColumnName(i), ((resultSet != null && resultSet.getObject(i) != null) ? resultSet.getObject(i) : "")); } catch (SQLException ignored) { ignored.printStackTrace(); } return null; }) .collect(Collectors.toMap(KV::getKey, stringObjectKV1 -> stringObjectKV1 != null ? stringObjectKV1.getValue() : ""));
问题原因
异常本质是Apache Beam的JdbcIO组件与Spanner JDBC驱动的事务处理逻辑冲突:
- Spanner JDBC默认开启autocommit,但执行查询时驱动会隐式开启事务;
- JdbcIO复用连接时,会尝试重置连接状态(包括设置autocommit),若上一个查询的事务未完全结束(比如ResultSet未正确关闭),就会触发该错误;
- 窗口处理带来的批量查询场景会提升连接复用概率,导致异常间歇性出现。
解决方案
方案1:禁用JdbcIO连接复用
在数据源配置中显式指定autocommit属性,并禁用连接池复用:
JdbcIO.<KV<String, Long>, KV<String, String>>readAll() .withDataSourceConfiguration( JdbcIO.DataSourceConfiguration.create( CONSTANTS.SPANNERDRIVER, "jdbc:cloudspanner:/projects/" + options.getInputProjectId() + "/instances/" + options.getInputInstanceId() + "/databases/" + options.getInputDatabaseId() ) .withConnectionProperties("autocommit=true;") .withMaxIdleConnections(0) // 禁用连接池,每次请求新建连接 )
方案2:确保ResultSet被正确关闭
在RowMapper中显式关闭ResultSet,避免残留事务:
.withRowMapper((JdbcIO.RowMapper<KV<String, String>>) resultSet -> { try { String key = resultSet.getString(CONSTANTS.ID); LOG.info("key value is: " + key); ResultSetMetaData metadata = resultSet.getMetaData(); // 构建row的原有逻辑... return KV.of(key, row.toString()); // 示例返回值,根据实际需求调整 } finally { try { resultSet.close(); } catch (SQLException e) { LOG.warning("关闭ResultSet失败: " + e.getMessage()); } } })
方案3:改用SpannerIO进行关联查询(推荐)
直接使用Beam专为Spanner设计的SpannerIO替代JdbcIO,避免JDBC层面的兼容性问题:
// 替换原JdbcIO部分代码 eData = changeStreamedData .apply("提取ID", Values.create()) .apply("Spanner关联查询", SpannerIO.read() .withSpannerConfig(spannerConfig) .withQuery(CONSTANTS.QUERY) .withParameterProvider(new SpannerIO.ParameterProvider<Long>() { @Override public StatementParameters getParameters(Long element) { return StatementParameters.newBuilder() .setLong(1, element) .build(); } }) .withRowMapper((Struct row) -> { String id = row.getString(CONSTANTS.ID); // 转换为KV格式的原有逻辑... return KV.of(id, row.toString()); }) );
总结
优先推荐方案3,SpannerIO能更好适配Beam的执行模型,从根源避免JDBC驱动带来的冲突。若必须使用JdbcIO,优先采用方案1禁用连接复用,同时配合方案2确保资源正确释放。
内容的提问来源于stack exchange,提问作者Nikhil
相关产品推荐
相关产品推荐

