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

Apache Beam DataFlow使用JdbcIO读取Oracle时触发空指针异常求助

Oracle 11g到BigQuery迁移时JdbcIO读取阶段的NullPointerException问题

使用Apache Beam Java SDK结合DataFlow Runner实现Oracle 11g到BigQuery的数据迁移,遇到以下特定场景的异常:

  • 当查询语句中选中列的总字节数超过~46752 Bytes时,JdbcIO.read()阶段触发NullPointerException
  • 总字节数低于该阈值时运行完全正常
  • 测试验证:列数超过300但总字节数未超阈值时可正常读取;空值、特定列组合均不是异常诱因

读取Oracle数据代码

// Read from JDBC
Pipeline p2 = Pipeline.create(
                PipelineOptionsFactory.fromArgs(args).withValidation().create());

String query2 = "SELECT Col1...Col00 FROM table WHERE rownum<=1000";
PCollection<TableRow> rows = p2.apply(JdbcIO.<TableRow>read()
                .withDataSourceConfiguration(JdbcIO.DataSourceConfiguration.create(
                                "oracle.jdbc.OracleDriver", "jdbc:oracle:thin:@//localhost:1521/orcl")
                        .withUsername("root")
                        .withPassword("password"))
                .withQuery(query2)
                .withRowMapper(new JdbcIO.RowMapper<TableRow>() {
                    @Override
                    public TableRow mapRow(ResultSet resultSet) throws Exception {

                        schema = getSchemaFromResultSet(resultSet);
                        TableRow tableRow = new TableRow();
                        List<TableFieldSchema> columnNames = schema.getFields();    
                        // 原代码省略字段填充逻辑
                        return tableRow;
                    }
                })
        );

p2.run().waitUntilFinish();

异常信息(仅列总字节数超阈值时触发)

Error message from worker: java.lang.NullPointerException
oracle.sql.converter.CharacterConverter1Byte.toUnicodeChars(CharacterConverter1Byte.java:344) 
oracle.sql.CharacterSet1Byte.toCharWithReplacement(CharacterSet1Byte.java:134)   
oracle.jdbc.driver.DBConversion._CHARBytesToJavaChars(DBConversion.java:964) 
oracle.jdbc.driver.DBConversion.CHARBytesToJavaChars(DBConversion.java:867)  
oracle.jdbc.driver.T4CVarcharAccessor.unmarshalOneRow(T4CVarcharAccessor.java:298)  
oracle.jdbc.driver.T4CTTIrxd.unmarshal(T4CTTIrxd.java:934)     
oracle.jdbc.driver.T4CTTIrxd.unmarshal(T4CTTIrxd.java:853)     
oracle.jdbc.driver.T4C8Oall.readRXD(T4C8Oall.java:699)     
oracle.jdbc.driver.T4CTTIfun.receive(T4CTTIfun.java:337)     
oracle.jdbc.driver.T4CTTIfun.doRPC(T4CTTIfun.java:191)     
oracle.jdbc.driver.T4C8Oall.doOALL(T4C8Oall.java:523)     
oracle.jdbc.driver.T4CPreparedStatement.doOall8(T4CPreparedStatement.java:207)     
oracle.jdbc.driver.T4CPreparedStatement.executeForDescribe(T4CPreparedStatement.java:863)     
oracle.jdbc.driver.OracleStatement.executeMaybeDescribe(OracleStatement.java:1153)     
oracle.jdbc.driver.OracleStatement.doExecuteWithTimeout(OracleStatement.java:1275)     
oracle.jdbc.driver.OraclePreparedStatement.executeInternal(OraclePreparedStatement.java:3576)     
oracle.jdbc.driver.OraclePreparedStatement.executeQuery(OraclePreparedStatement.java:3620)     
oracle.jdbc.driver.OraclePreparedStatementWrapper.executeQuery(OraclePreparedStatementWrapper.java:1491)     
org.apache.commons.dbcp2.DelegatingPreparedStatement.executeQuery(DelegatingPreparedStatement.java:122)     
org.apache.commons.dbcp2.DelegatingPreparedStatement.executeQuery(DelegatingPreparedStatement.java:122)     
org.apache.beam.sdk.io.jdbc.JdbcIO$ReadFn.processElement(JdbcIO.java:1381)

排查方向与解决建议

1. 校验Oracle JDBC驱动版本

异常发生在Oracle JDBC内部的字符转换逻辑,大概率是驱动版本与Oracle 11g不兼容或驱动存在已知bug:

  • 更换为适配Oracle 11g的稳定驱动版本(推荐ojdbc6或ojdbc7,避免使用ojdbc8及以上版本,这类高版本对11g的兼容性较差)
  • 替换驱动后重新执行任务测试

2. 调整JdbcIO的连接与读取参数

通过调整连接属性减少单批次读取的数据量,降低内存压力:

  • 在数据源配置中添加defaultRowPrefetch属性,控制单次从数据库读取的行数:
.withDataSourceConfiguration(JdbcIO.DataSourceConfiguration.create(
        "oracle.jdbc.OracleDriver", "jdbc:oracle:thin:@//localhost:1521/orcl")
        .withUsername("root")
        .withPassword("password")
        .withConnectionProperties("defaultRowPrefetch=100"))
  • 若使用的Beam版本支持,可直接为JdbcIO设置fetchSize:.withFetchSize(100)

3. 统一字符集配置

异常出现在字符转换环节,需确保JDBC连接字符集与Oracle数据库字符集一致:

  • 在JDBC URL中明确指定字符集参数,例如数据库使用UTF8时:jdbc:oracle:thin:@//localhost:1521/orcl?characterEncoding=UTF-8

4. 拆分查询绕过字节限制

如果上述方法均无效,可临时拆分查询:

  • 将大字段或多字段拆分为多个独立的读取任务,在Beam中通过键值关联合并TableRow,规避单条记录总字节数的限制

内容的提问来源于stack exchange,提问作者NIKHIL SUTHAR

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 18:01:23