Dataproc上Pyspark处理GCS日志实现BigQuery增量写入方案咨询
Pyspark逐GCS文件写入BigQuery实现方案
核心实现逻辑
- 先通过Hadoop FileSystem API拉取GCS指定目录下所有
.csv.gz格式的文件完整路径列表 - 遍历每个文件路径,单独读取为DataFrame后写入BigQuery
- 首次写入采用
overwrite模式清空上周旧数据,后续所有文件采用append模式追加写入,每写完一个文件即可在BigQuery中查询到对应部分的数据
修改后可运行代码
import pyspark from pyspark.sql.types import StructType, StructField, StringType, DateType from pyspark.sql import SparkSession spark = SparkSession.builder.appName("Node Log ETL").getOrCreate() sc = spark.sparkContext # -------------------------- 配置部分 -------------------------- # 日志文件存放的GCS目录 LOG_GCS_DIR = "gs://some-bucket/" # BigQuery写入配置 BQ_TEMP_BUCKET = "some-bucket/node-logs-temp-bq" BQ_DATASET = "somedataset" BQ_TABLE = "node_weekly_logs" # 预先定义日志Schema,替换为你实际的字段结构 LOG_SCHEMA = StructType([ StructField("log_id", StringType(), nullable=True), StructField("node_ip", StringType(), nullable=True), StructField("event_date", DateType(), nullable=True), # 按需补充其余所有字段 ]) # -------------------------- 功能实现 -------------------------- # 调用Hadoop API获取GCS目录下所有csv.gz文件路径 URI = sc._gateway.jvm.java.net.URI Path = sc._gateway.jvm.org.apache.hadoop.fs.Path FileSystem = sc._gateway.jvm.org.apache.hadoop.fs.FileSystem Configuration = sc._gateway.jvm.org.apache.hadoop.conf.Configuration fs = FileSystem.get(URI(LOG_GCS_DIR), Configuration()) file_status_list = fs.listStatus(Path(LOG_GCS_DIR)) # 筛选出所有.csv.gz后缀的文件路径 csv_gz_files = [] for file_status in file_status_list: file_full_path = file_status.getPath().toString() if file_full_path.endswith(".csv.gz"): csv_gz_files.append(file_full_path) # 逐文件读取并写入BigQuery is_first_write = True for idx, file_path in enumerate(csv_gz_files, 1): # 读取单个csv.gz文件 log_df = spark.read.csv( path=file_path, header=True, schema=LOG_SCHEMA ) # 写入BQ,首次覆盖旧表,后续追加 write_mode = "overwrite" if is_first_write else "append" log_df.write.format("bigquery")\ .option("temporaryGcsBucket", BQ_TEMP_BUCKET)\ .option("dataset", BQ_DATASET)\ .mode(write_mode)\ .save(BQ_TABLE) is_first_write = False print(f"进度:{idx}/{len(csv_gz_files)},文件{file_path}已写入BigQuery,可直接查询")
注意事项
- 请提前将代码中的
LOG_SCHEMA替换为你实际日志文件的字段结构,预定义Schema可以避免Spark逐文件推断Schema的性能损耗,也能避免不同文件字段不一致导致的写入失败 - 如果需要保留多周历史日志,建议给BigQuery表增加周数分区字段,全程使用
append模式写入即可,无需首次覆盖 - 可根据需要增加异常捕获逻辑,记录写入失败的文件路径,后续单独重跑失败文件即可,无需重跑全量
- gzip为不可拆分压缩格式,单个gz文件默认由1个Executor处理,符合逐文件处理的预期,不会影响原有集群的并行处理能力
内容的提问来源于stack exchange,提问作者Timothy
相关产品推荐
相关产品推荐

