PySpark循环遍历日期读取数据湖数据并调用独立Notebook处理
实现日期循环中调用Notebook完成数据湖Parquet数据处理流程
问题背景
数据湖内的Parquet数据按yyyy/mm/dd目录结构化存储,需求是对每个目标日期:
- 拉取该日期最近5天的原始数据
- 执行预设的数据转换逻辑
- 仅保留当前目标日期的记录
- 将结果写入数据湖对应日期的目录
现有日期循环逻辑,但需要补充在循环中调用独立Notebook完成完整数据处理流程的实现。
步骤1:修正日期范围逻辑
原代码中计算的日期范围仅覆盖4天数据,不符合“最近5天”的需求,先修正这部分:
from datetime import datetime, timedelta first_date = datetime(2022, 2, 20).date() last_date = datetime(2022, 2, 12).date() current_date = first_date while current_date >= last_date: print(f"处理目标日期:{current_date}") # 计算最近5天的日期范围:current_date 到 current_date-4(共5天) file_paths = [] for offset in range(5): target_date = current_date - timedelta(days=offset) file_path = f"/mnt/container/{target_date.strftime('%Y/%m/%d')}/*.parquet" file_paths.append(file_path) print(f"待读取的文件路径:{file_paths}") # --- 核心:调用独立Notebook处理数据 --- notebook_params = { "input_paths": ",".join(file_paths), # 用逗号拼接路径,方便Notebook解析 "target_date": current_date.strftime('%Y-%m-%d'), "output_base_path": "/mnt/container/processed" # 处理后数据的根目录 } # 调用独立Notebook(以Databricks环境为例) dbutils.notebook.run( path="/path/to/your/transformation_notebook", # 替换为你的Notebook实际路径 timeout_seconds=3600, arguments=notebook_params ) current_date = current_date - timedelta(days=1)
步骤2:实现独立转换Notebook的逻辑
在指定路径的Notebook中,接收参数并完成数据读取、转换、过滤、写入:
# 1. 接收外部传入的参数 dbutils.widgets.text("input_paths", "") dbutils.widgets.text("target_date", "") dbutils.widgets.text("output_base_path", "") input_paths = dbutils.widgets.get("input_paths").split(",") target_date = dbutils.widgets.get("target_date") output_base_path = dbutils.widgets.get("output_base_path") # 2. 读取指定路径的Parquet数据 df = spark.read.parquet(*input_paths) # 3. 执行预设的数据转换逻辑(替换为你的实际转换代码) # 示例:添加计算列、清洗脏数据等 transformed_df = df.withColumn("processed_time", current_timestamp()) # 4. 过滤仅保留当前目标日期的记录 # 假设原始数据中有`date`字段存储日期,格式为'yyyy-MM-dd' from pyspark.sql.functions import col filtered_df = transformed_df.filter(col("date") == target_date) # 5. 写入数据湖对应日期目录 output_path = f"{output_base_path}/{target_date.replace('-', '/')}" # 覆盖模式:若目录已存在则替换,根据需求可改为append filtered_df.write.mode("overwrite").parquet(output_path) # 清理临时widgets(可选) dbutils.widgets.removeAll()
关键说明
- 参数传递:通过逗号拼接路径列表,避免复杂的参数格式解析;日期统一用
yyyy-MM-dd字符串传递,避免跨Notebook的日期类型序列化问题 - 写入模式:使用
overwrite模式确保每次处理的结果是最新的,若需要增量追加可改为append - 超时设置:根据数据量调整
timeout_seconds,避免处理大文件时超时
内容的提问来源于stack exchange,提问作者asd
相关产品推荐
相关产品推荐

