【发布时间】:2021-07-04 16:57:35
【问题描述】:
我正在尝试在这个 google public BigQuery table 上使用 pySpark(表大小:268.42 GB,行数:611,647,042)。我将集群的区域设置为 US(与 BigQuery 表相同),但即使在集群中使用多台高性能机器时,代码也非常慢。知道为什么吗?我应该在我的存储桶中创建公共 BigQuery 表的副本吗?如果是,怎么做?
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.master('yarn') \
.appName('spark-bigquery-crypto') \
.config('spark.jars', 'gs://spark-lib/bigquery/spark-bigquery-latest_2.12.jar') \
.getOrCreate()
# Use the Cloud Storage bucket for temporary BigQuery export data used
# by the spark-bigquery-connector.
bucket = "dataproc-staging-us-central1-397704471406-lrrymuq9"
spark.conf.set('temporaryGcsBucket', bucket)
# Load data from BigQuery.
eth_transactions = spark.read.format('bigquery') \
.option('table', 'bigquery-public-data:crypto_ethereum.transactions') \
.load()
eth_transactions.createOrReplaceTempView('eth_transactions')
# Perform SQL query.
df = spark.sql('''SELECT * FROM eth_transactions WHERE DATE(block_timestamp) between "2019-01-01" and "2019-01-31"''')
【问题讨论】:
标签: google-cloud-platform pyspark google-bigquery dataproc