Apache Spark对接BigQuery时如何执行自定义DDL/DML类型SQL查询
核心结论
Spark BigQuery 连接器的spark.read()接口是专门用于读取数据的,不支持直接传入多语句DDL+DML混合SQL执行,不能像JDBC那样直接把整段SQL填入table参数运行。你可以选择以下两种方案实现需求:
方案一:先通过BigQuery Java客户端执行自定义SQL,再用Spark读取结果表
该方案可以直接复用你已有的SQL逻辑,改动最小:
- 先引入BigQuery官方Java客户端依赖
- 用客户端提交多语句SQL,等待执行完成后,再用Spark读取生成好的目标表即可
示例代码:
// 初始化BigQuery客户端 BigQuery bigquery = BigQueryOptions.newBuilder() .setProjectId("customerProject") .setCredentials(ServiceAccountCredentials.fromStream(new ByteArrayInputStream(serviceAccountKey.getBytes()))) .build() .getService(); // 你的自定义业务SQL String customSql = "DROP TABLE DEMO.CUSTOMER;\n" + "CREATE TABLE DEMO.CUSTOMER AS (\n" + "SELECT DISTINCT\n" + " CUSTOMERID,\n" + " CUSTOMERNAME\n" + " FROM DEMO.COMPANY);\n" + "INSERT INTO DEMO.CUSTOMER SELECT CUSTOMERID,CUSTOMERNAME FROM DEMO.COMPANY WHERE AV = 1;"; // 配置并提交SQL作业,开启多语句执行支持 QueryJobConfiguration queryConfig = QueryJobConfiguration.newBuilder(customSql) .setAllowLargeResults(true) .setMultiStatementAllowed(true) .build(); Job job = bigquery.create(JobInfo.of(queryConfig)); job = job.waitFor(); // 执行成功后再用Spark读取目标表 if (job.isDone() && job.getStatus().getError() == null) { HashMap<String, String> options = new HashMap<>(); options.put("credentials", serviceAccountKey); options.put("project", "customerProject"); options.put("parentProject", "customerProject"); options.put("table", "DEMO.CUSTOMER"); Dataset<Row> result = spark.read().format("bigquery").options(options).load(); } else { throw new RuntimeException("SQL执行失败:" + job.getStatus().getError().toString()); }
方案二:完全基于Spark实现相同逻辑,无需原生BigQuery SQL
如果不想额外引入BigQuery客户端,可以把SQL逻辑转换为Spark算子实现,overwrite写入模式会自动覆盖原有表,等价于删表重建的效果:
// 读取源表配置 HashMap<String, String> readOptions = new HashMap<>(); readOptions.put("credentials", serviceAccountKey); readOptions.put("project", "customerProject"); readOptions.put("parentProject", "customerProject"); readOptions.put("table", "DEMO.COMPANY"); Dataset<Row> sourceData = spark.read().format("bigquery").options(readOptions).load(); // 实现原SQL的去重+过滤合并逻辑 Dataset<Row> distinctData = sourceData.select("CUSTOMERID", "CUSTOMERNAME").distinct(); Dataset<Row> filteredData = sourceData.filter(col("AV").equalTo(1)).select("CUSTOMERID", "CUSTOMERNAME"); Dataset<Row> resultData = distinctData.union(filteredData); // 写入目标表 HashMap<String, String> writeOptions = new HashMap<>(); writeOptions.put("credentials", serviceAccountKey); writeOptions.put("project", "customerProject"); writeOptions.put("parentProject", "customerProject"); writeOptions.put("table", "DEMO.CUSTOMER"); resultData.write().format("bigquery").options(writeOptions).mode("overwrite").save(); // 后续需要读取直接访问目标表即可 Dataset<Row> customerData = spark.read().format("bigquery").options(writeOptions).load();
两种方案按需选择即可:方案一适合SQL逻辑复杂、不想改写的场景;方案二适合逻辑简单、希望全链路用Spark调度的场景。
内容的提问来源于stack exchange,提问作者David Wallace
相关产品推荐
相关产品推荐

