Dataproc中使用PySpark读写BigQuery数据时的错误排查
问题:PySpark通过SQL读取BigQuery数据时触发'dataset未提供'错误
问题背景
在Dataproc Workbench的用户托管Jupyter Notebook实例中,使用PySpark读写BigQuery数据(目标表ID:my-project.mydatabase.mytable)。已添加BQ连接器,全表读取代码可正常运行,但会导入全量1.09亿条记录、9列数据,实际仅需100万条记录的posting列。尝试用SQL查询读取时触发错误。
错误信息
Py4JJavaError: An error occurred while calling o195.load. : com.google.cloud.spark.bigquery.repackaged.com.google.inject.ProvisionException: Unable to provision, see the following errors: 1) Error in custom provider, java.lang.IllegalArgumentException: 'dataset' not parsed or provided. at com.google.cloud.spark.bigquery.SparkBigQueryConnectorModule.provideSparkBigQueryConfig(SparkBigQueryConnectorModule.java:65) while locating com.google.cloud.spark.bigquery.SparkBigQueryConfig 1 error at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.InternalProvisionException.toProvisionException(InternalProvisionException.java:226) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.InjectorImpl$1.get(InjectorImpl.java:1097) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.InjectorImpl.getInstance(InjectorImpl.java:1131) at com.google.cloud.spark.bigquery.BigQueryRelationProvider.createRelationInternal(BigQueryRelationProvider.scala:75) at com.google.cloud.spark.bigquery.BigQueryRelationProvider.createRelation(BigQueryRelationProvider.scala:46) at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:332) at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:242) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:230) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:197) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayConnection.run(GatewayConnection.java:238) at java.lang.Thread.run(Thread.java:750) Caused by: java.lang.IllegalArgumentException: 'dataset' not parsed or provided. at com.google.cloud.bigquery.connector.common.BigQueryUtil.lambda$parseTableId$2(BigQueryUtil.java:153) at java.util.Optional.orElseThrow(Optional.java:290) at com.google.cloud.bigquery.connector.common.BigQueryUtil.parseTableId(BigQueryUtil.java:153) at com.google.cloud.spark.bigquery.SparkBigQueryConfig.from(SparkBigQueryConfig.java:237) at com.google.cloud.spark.bigquery.SparkBigQueryConnectorModule.provideSparkBigQueryConfig(SparkBigQueryConnectorModule.java:67) at com.google.cloud.spark.bigquery.SparkBigQueryConnectorModule$$FastClassByGuice$$db983008.invoke(<generated>) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.ProviderMethod$FastClassProviderMethod.doProvision(ProviderMethod.java:264) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.ProviderMethod.doProvision(ProviderMethod.java:173) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.InternalProviderInstanceBindingImpl$CyclicFactory.provision(InternalProviderInstanceBindingImpl.java:185) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.InternalProviderInstanceBindingImpl$CyclicFactory.get(InternalProviderInstanceBindingImpl.java:162) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.ProviderToInternalFactoryAdapter.get(ProviderToInternalFactoryAdapter.java:40) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.SingletonScope$1.get(SingletonScope.java:168) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.InternalFactoryToProviderAdapter.get(InternalFactoryToProviderAdapter.java:39) at com.google.cloud.spark.bigquery.repackaged.com.google.inject.internal.InjectorImpl$1.get(InjectorImpl.java:1094) ... 18 more
用户代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import udf, col from pyspark.sql.types import IntegerType, ArrayType, StringType from google.cloud import bigquery # 添加BQ连接器 spark = SparkSession.builder.appName('SpacyOverPySpark') \ .config('spark.jars.packages', 'com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.24.2') \ .getOrCreate() # 全表读取(可正常运行) df = spark.read.format('bigquery').option('table', 'my-project.mydatabase.mytable').load() print("DataFrame shape: ", (df.count(), len(df.columns))) # 输出109M记录 & 9列 # 失败的SQL读取尝试 # sql = """ # SELECT `posting` # FROM `mentor-pilot-project.indeed.indeed-data-clean` # LIMIT 1000000 # """ # df = spark.read.format("bigquery").load(sql) # print("DataFrame shape: ", (df.count(), len(df.columns))) # 应急方案:从Cloud Storage读取(可正常运行) df = spark.read.csv("gs://hidden_bucket/1M_samples.csv", header=True) # 示例自定义处理(可正常运行) def split_text(text:str) -> list: return text.split() textsplitUDF = udf(lambda z: split_text(z), ArrayType(StringType())) df.withColumn("posting_split", textsplitUDF(col("posting"))) # 写入BigQuery的正确代码(已解决写入问题) df.write \ .format('bigquery') \ .option("temporaryGcsBucket", "my_temp_bucket_name") \ .mode("overwrite") \ .save("my-project.mynewdatabase.mytable")
问题原因与解决方案
错误原因
spark.read.format("bigquery").load(sql)的用法不符合BigQuery连接器规范:load()方法仅接收表路径或标识符,不能直接传入SQL查询语句。连接器会把传入的SQL当成BigQuery表ID解析,无法识别出有效dataset,因此触发'dataset' not parsed or provided错误。
正确方案
有两种可行方式实现SQL查询读取BigQuery数据:
方案1:使用query参数指定SQL语句
将SQL查询通过option("query", sql)传递给读取器,而非传给load():
sql = """ SELECT `posting` FROM `mentor-pilot-project.indeed.indeed-data-clean` LIMIT 1000000 """ df = spark.read.format("bigquery").option("query", sql).load() print("DataFrame shape: ", (df.count(), len(df.columns)))
方案2:创建临时视图后用Spark SQL查询
先将BigQuery表注册为Spark临时视图,再通过Spark SQL执行查询:
# 注册BigQuery表为临时视图 spark.read.format('bigquery').option('table', 'mentor-pilot-project.indeed.indeed-data-clean').createOrReplaceTempView('bq_table') # 执行Spark SQL查询 df = spark.sql("SELECT `posting` FROM bq_table LIMIT 1000000") print("DataFrame shape: ", (df.count(), len(df.columns)))
内容的提问来源于stack exchange,提问作者David Espinosa
相关产品推荐
相关产品推荐

