You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.30 09:25:00