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
相关产品推荐
相关产品推荐

