Google Cloud Dataflow/Apache Beam如何将CloudSQL取到的LAST_UPD_TS设为BQ查询参数
问题原因
- 核心逻辑错误:你混淆了Apache Beam流水线的构建阶段和运行阶段的边界。IO转换(比如BigQueryIO、JdbcIO)属于流水线的静态拓扑定义,必须在
pipeline.run()执行前的构建阶段声明,不能在运行时的ParDo处理函数中创建IO转换实例并作为数据输出,你后续直接给JdbcIO.write传入IO转换类型的PCollection也完全不符合输入要求,自然无法正常执行。 - 空指针异常的直接诱因:你调用了
.fromQuery()方法指定BigQuery的查询来源,而非.fromTable()指定固定表,此时getTable()方法返回值为null,调用.get("loc_dim")必然抛出空指针。
正确实现方案
你需要用BigQueryIO提供的readAll()方法实现动态查询:该方法支持接收一个包含查询语句的PCollection作为输入,运行时会动态执行PCollection中的所有查询,返回查询结果的TableRow集合,刚好匹配你的需求。
// 第一步:从CloudSQL读取最大LAST_UPD_TS PCollection<String> maxTsCollection = pipeline.apply("get latest ts", JdbcIO.<String>read() .withDataSourceConfiguration(JdbcIO.DataSourceConfiguration.create( ValueProvider.StaticValueProvider.of("com.mysql.jdbc.Driver"), jdbcUrlValueProvider)) .withQuery("select MAX(last_upd_ts) AS last_upd_ts from DEPT") .withCoder(StringUtf8Coder.of()) // 直接用字符串编码器更合适,不需要AvroCoder .withRowMapper(resultSet -> { String maxTs = resultSet.getString("last_upd_ts"); // 处理表为空返回null的情况,可替换为你业务需要的最小时间 return maxTs == null ? "1970-01-01 00:00:00" : maxTs; }) ); // 第二步:将最大TS拼接为BQ查询语句 PCollection<String> bqQueryCollection = maxTsCollection.apply("Build BQ Query", ParDo.of(new DoFn<String, String>() { @ProcessElement public void processElement(@Element String maxTs, ProcessContext c) { // 注意时间类型需要加单引号包裹,避免SQL语法错误 String query = String.format("select * from `pdata.DEPT` WHERE LAST_UPD_TS >= '%s'", maxTs); c.output(query); } })); // 第三步:动态执行BQ查询,拿到结果 PCollection<TableRow> bqResult = bqQueryCollection.apply("Read from BQ dynamically", BigQueryIO.readTableRows() .withTemplateCompatibility() .withoutValidation() .usingStandardSql() .readAll() // 核心方法:接收PCollection中的查询语句动态执行 ); // 第四步:将结果写入CloudSQL,复用你原来的逻辑即可 bqResult.apply("Insert and Update", JdbcIO.<TableRow>write() .withDataSourceConfiguration(JdbcIO.DataSourceConfiguration.create( ValueProvider.StaticValueProvider.of("com.mysql.jdbc.Driver"), jdbcUrlValueProvider)) .withStatement("insert into DEPT (LOC_DIM_ID,DIVN_NBR,DEPT_NBR,END_DT,START_DT,PRC_OPT_CD,PRN_LVL_CD,PRICE_LOC_NBR,LAST_UPD_TS,LAST_UPD_USERID)" + "values( ?,?,?,?,?,?,?,?,?,?)" + "ON DUPLICATE KEY UPDATE START_DT=?,PRC_OPT_CD=?,PRN_LVL_CD=?,PRICE_LOC_NBR=?,LAST_UPD_TS=?,LAST_UPD_USERID=?") .withPreparedStatementSetter(new DEPT_BULKPreparedStatementSetters()) ); PipelineResult.State state = pipeline.run().waitUntilFinish();
内容的提问来源于stack exchange,提问作者Akshay
相关产品推荐
相关产品推荐

