Databricks上PySpark读取BigQuery分区必填表报错及依赖冲突问题
解决方案:Databricks PySpark读取BigQuery强制分区表及多数据源冲突问题
问题1:强制分区过滤表读取失败(方法1)
BigQuery强制分区过滤要求必须在查询阶段就在BigQuery端应用分区条件,而Spark DataFrame的filter()/where()是数据加载到Spark后才执行的过滤,无法满足BigQuery的分区强制要求。需要通过BigQuery Connector的专属参数将过滤条件传递到BigQuery端执行,以下是两种可行方式:
方式A:使用query参数直接执行SQL
直接通过query选项传入包含分区过滤的完整SQL,让BigQuery先执行过滤后再返回数据:
read_df = spark.read.format('bigquery') \ .option("parentProject","valid_parent_project_name")\ .option("project", "valid_project_name")\ .option("dataset", "valid_dataset_name")\ .option('query', "SELECT * FROM investment WHERE TIMESTAMP_TRUNC(update_date, DAY) = TIMESTAMP('2023-10-10')") \ .load()
方式B:使用filter选项传递分区条件
通过Connector的filter参数(注意不是DataFrame的filter方法)直接指定分区过滤规则,BigQuery会自动应用该条件:
read_df = spark.read.format('bigquery') \ .option("parentProject","valid_parent_project_name")\ .option("project", "valid_project_name")\ .option("dataset", "valid_dataset_name")\ .option('table', "investment") \ .option("filter", "TIMESTAMP_TRUNC(update_date, DAY) = TIMESTAMP('2023-10-10')") \ .load()
问题2:多BigQuery数据源冲突(方法2)
报错是因为Databricks自带的BigQuery Connector与你手动添加的com.google.cloud.spark:spark-3.3-bigquery:0.32.2库存在Provider冲突,解决方法是指定完全限定的Provider类名,同时修正配置参数的错误:
修正后的代码
# 配置临时物化表的项目和数据集(注意不是表名) spark.conf.set("viewsEnabled","true") spark.conf.set("materializationDataset","valid_project_name.valid_dataset_name") spark.conf.set("materializationProject", "valid_parent_project_name") # 指定完整Provider类名,并通过query参数传递SQL sql = "SELECT * FROM valid_project_name.valid_dataset_name.investment WHERE TIMESTAMP_TRUNC(update_date, DAY) = TIMESTAMP('2023-10-10') LIMIT 1000" df = spark.read.format("com.google.cloud.spark.bigquery.v2.Spark33BigQueryTableProvider") \ .option("query", sql) \ .load() display(df)
关键修正点
materializationDataset和materializationProject只需指定项目和数据集,不需要加表名;- 必须在
format()中写入完整的Provider类名,二选一即可:- 旧版Provider:
com.google.cloud.spark.bigquery.BigQueryRelationProvider - 新版Spark3.3适配Provider:
com.google.cloud.spark.bigquery.v2.Spark33BigQueryTableProvider
- 旧版Provider:
- 不要直接将SQL传入
load()方法,需通过option("query", sql)传递查询语句。
内容的提问来源于stack exchange,提问作者Manigandan
相关产品推荐
相关产品推荐

