You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

从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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 06:37:04