← 返回蜂巢洞察

ETL流程构建手册:如何使用Python搭建一套可用于生产环境的高效流程系统

要追踪洪水风险,需要一样虽然不起眼但却至关重要的东西:干净且结构清晰的数据。 在本教程中,你将亲手构建一个数据处理流程。你会创建一个Python ETL(提取、转换、加载)流程,该流程会从法国官方开放的水文数据API Hub'Eau 中获取每日的水位数据。然后,你需要对这些数据进行清洗,并将其发布为公共数据集,就像 实时版本 所做的那样。 本教程基于一个实际存在的流程系统——该系统每周会按计划自动运行一次,从而确保 巴黎洪水数据集 能够得到及时更新。 不过,你并不会只是简单地复制和粘贴代码。真正的目标是理解这个流程为何能以这样的方式运作。你会了解到,那些让一个脚本能够仅运行一次,却能让另一个脚

要追踪洪水风险,需要一样虽然不起眼但却至关重要的东西:干净且结构清晰的数据。

在本教程中,你将亲手构建一个数据处理流程。你会创建一个Python ETL(提取、转换、加载)流程,该流程会从法国官方开放的水文数据API Hub'Eau中获取每日的水位数据。然后,你需要对这些数据进行清洗,并将其发布为公共数据集,就像实时版本所做的那样。

本教程基于一个实际存在的流程系统——该系统每周会按计划自动运行一次,从而确保巴黎洪水数据集能够得到及时更新。

不过,你并不会只是简单地复制和粘贴代码。真正的目标是理解这个流程为何能以这样的方式运作。你会了解到,那些让一个脚本能够仅运行一次,却能让另一个脚本多年持续无人值守地正常运行的设计决策究竟是什么。

你可以配合这份笔记本来学习编程。在教程的大部分内容中,流程都是使用模拟API数据进行运行的。这样你就可以安全地尝试每一行代码,而不会对真实的服务器造成负担。后续的部分会介绍如何切换到实时API。

学完本教程后,你将能够:

  • 理解并掌握提取、转换、加载这三种数据处理步骤的原理

  • 使用Python的`@dataclass`语法来管理配置信息,而不是把各种常量分散地写在代码各处

  • 编写能够在网络故障或数据分页传输的情况下正常运行的API调用代码

  • 采用可靠的类型转换机制,确保某个错误数据不会导致整个流程崩溃

  • 安全地删除重复数据,并合并增量式更新的数据

  • 将所有组件连接成一个单一的、可调度执行的`main()`函数

目录:

先决条件

请安装以下依赖库:

pip install requests pandas numpy ipykernel

你将在Jupyter笔记本中完成相关操作。可以从Kaggle下载.ipynb文件,或者直接点击“复制并编辑”在Kaggle平台上进行操作。为此,你需要拥有一个Kaggle账户。

用于在Kaggle上下载数据工程学习笔记本的菜单按钮。

可选:如果你打算在本地运行这个笔记本,还需要安装notebook包:

pip install notebook

第1部分:设计逻辑

在开始编写任何代码之前,我们先花三分钟时间来探讨一下为什么这个数据管道要被设计成这样的结构。

这就是整个项目的整体架构。后续每一个代码层面的决策,其实都可以追溯到这些初始设计理念之中。即使你跳过了其他内容,也请一定要阅读这一部分。

什么是ETL管道?

ETL代表提取、转换、加载。这是一种将数据从源位置可靠且可重复地传输到目标位置的标准流程。

  • 提取:从数据库、API或文件等来源中获取数据。

  • 转换:对数据进行清洗、标准化、补充信息或验证其准确性。

  • 加载:将处理后的数据写入目标存储位置,如数据仓库、CSV文件或公共平台。

下面来看这三个阶段是如何应用到这个项目中的:

阶段 具体操作内容
提取 如果现有数据存在,则先将其加载进来;然后针对每个测量站,通过requests调用Hub'Eau API来获取数据。
转换 将法语列和条目的内容转换为英语,调整数据类型,并删除重复项。
加载 生成CSV文件和元数据文件,然后通过Kaggle CLI将其发布到平台上。

架构概览

简单ETL流程图。

请注意,“提取”阶段已经涉及两种不同的数据源:现有的数据集,以及从API获取的新数据。这种区分正是接下来要讨论的内容的基础。

使其具备生产级特性的两种模式

有这两种模式的存在,使得这个处理流程能够被安全地设置为定期自动运行,且可以持续运行多年而不会出现问题。

1. 等幂性

等幂性是指当同一操作被执行两次时,其结果与只执行一次时的结果相同。去重机制正是使这个处理流程具备等幂性的关键因素。如果调度系统意外地多次触发该操作,或者网络请求在当天再次尝试获取数据,第二次执行也不会产生重复的记录。对于那些需要定期自动运行且无人监控的系统来说,这一点非常重要。

2. 增量加载

一个简单的处理流程会在每次执行时重新下载所有数据。这样做不仅速度慢,还会浪费API的使用额度,而且系统也相当脆弱:传输的数据量越大,出现故障的概率也就越高。

增量加载机制可以避免这些问题。具体来说,它会:

  • 首先检查数据集中最新的日期

  • 仅请求从该日期之后的数据

  • 将新获取的记录合并到现有的数据集中

你可以在determine_update_range中看到这种逻辑的具体实现。

永远要问自己:这个习惯能够区分一个脆弱的脚本和一个具备生产级稳定性的处理流程。例如,在这个处理流程中,每一个“耗时较多”或“依赖外部资源”的操作(如网络请求、磁盘写入、数据发布等),都会先经过简单的本地检查才能继续执行。

第2部分:配置与依赖关系

以下是导入库的代码段。虽然看起来并不复杂,但其组织方式本身就属于值得推崇的最佳实践。

# 标准库导入
import json
import os
import random
import subprocess
from dataclasses import dataclass, field
from datetime import date, timedelta
from pathlib import Path
from typing import Dict, List, Optional, Set, Tuple

# 第三方库导入
import numpy as np
import pandas as pd
import requests

# 为Jupyter笔记本设置显示选项
pd.set_option('display.max_columns', None)   # 显示所有列
pd.set_option('display.max_colwidth', None)  # 防止长文本被截断

代码层面的说明:

  • 导入语句分为三类:标准库模块、第三方库模块,以及本地模块,各类导入语句之间用空行分隔。这种排列方式遵循了PEP 8关于导入顺序的规定。大多数自动格式化工具(如isortruff)也会采用同样的分类规则。

  • 切勿使用from module import *这种语法进行导入,因为这样会污染当前的命名空间,导致六个月后别人调试代码时难以追踪某个变量的来源。Python的各种风格指南也都建议避免使用这种写法,包括Google的Python风格指南

  • 函数pd.set_option(...)的存在纯粹是为了提高代码的可读性,这样当数据框的尺寸过大时,使用这些选项也不会导致数据被截断。这些选项对数据处理流程的逻辑本身没有任何影响;在普通的`.py`脚本中,通常应该删除这些代码或调整它们的作用范围。

第3部分:使用数据类管理配置

为什么要专门设置配置层呢?

任何数据处理流程都包含一些配置参数,比如需要监控哪些站点、洪水阈值的设定值,以及数据应该发布到哪里。如果将这些参数直接以字面量的形式分散写在代码中,就会导致代码结构变得混乱。例如,某个条件判断语句可能会深藏在函数结构的第三层之中,像if level > 6000:这样的代码。一旦需要修改其中任何一项配置,就不得不在整个文件中逐行查找相关内容,很容易遗漏某些地方。

解决办法是:将所有配置参数集中放在一个地方。Python中的@dataclass装饰器正是实现这一目标的最佳工具。

功能 重要性
自动生成的__init____repr____eq__方法 无需自行编写这些方法
类型提示功能 可帮助IDE自动补全代码,并使代码具有自文档化功能
可选的frozen=True选项 如果希望配置参数在创建后不可更改,该选项能确保其具备真正的不可变性
__post_init__钩子方法 在对象构造完成后,会一次性验证或计算一些派生字段的值

将这种写法与使用普通dict结构来存储配置进行对比:

config = {
    "stations": ["STN001", "STN002"],
    "flood_threshold": 6000,
    "publish_url": "https://example.com/alerts",
    "retry_count": 2,
    "timeout_seconds": 5,
}
<一个普通的字典根本无法提供这些功能:既没有不可变性,也没有类型检查机制,更没有自动完成输入的功能。>

代码级讲解

该处理流程定义了三个配置类。每个类都承担一项特定的职责:API详细信息、站点规则以及数据发布目标。每个类都会被创建成一个模块级别的单例实例,流程中的其他所有函数都会从这些单例实例中获取所需数据。

@dataclass
class APIConfig:
    """用于HubEau数据获取的API配置信息。

    可将这个类视为API的“地址簿”。”
    """
    use mocks: bool = True
    base_url: str = "https://hubeau.eaufrance.fr/api/v2/hydrometrie/obs_elab"
    metric: str = "HIXnJ"  # 日最高水位(详细观测数据)
    # 分页限制:每次请求最多获取20,000条记录
    max_per_page: int = 20000
    timeout_seconds: int = 60  # 网络超时时间

    def __post_init__(self):
        """在初始化后验证配置信息是否正确。”
        if self.max_per_page <= 0:
            raise ValueError("max_per_page必须为正数")
        if self.timeout_seconds <= 0:
            raise ValueError("timeout_seconds必须为正数")

请注意__post_init__方法中的验证逻辑。这个验证会在自动生成的__init__方法执行后立即进行。如果APIConfig(max_per_page=-1)这样的错误配置被使用,程序在启动时会立即出现异常,而不会在实际的数据处理流程中才出现问题。

@dataclass
class StationConfig:
    """站点监控配置信息。

    这决定了数据收集的具体内容:哪些站点需要被监测,如何判断是否发生洪水?”
    """
    station_codes: List[str] = field(default_factory=lambda: [
        "F700000109", "F700000110", "F700000111",
        "F700000102", "F700000103",
    ])
    flood_threshold_mm: int = 6000  # 洪水警报阈值
    earliest_date: str = "1900-01-01"  # 数据追溯的起始日期

关于可变默认值的注意事项:请仔细观察station_codes的定义方式。它并没有被写成station_codes: List[str] = [...]这种形式,这是有意为之,因为这样可以避免Python中一个常见的错误。如果使用可变对象(如列表、字典或集合)作为默认参数或数据类字段,所有实例都会共享同一个对象。在某个实例上对该对象进行修改,其他所有实例也会随之被修改。

Stack Overflow在其关于“最小惊讶原则”与可变默认参数的讨论中详细介绍了这个问题,Real Python关于可选参数的指南中也有所涉及。解决这个问题的方法是使用field(default_factory=...)这种写法。这样,每次创建新实例时,都会调用一个新的工厂函数(这里是lambda表达式),从而确保每个实例都拥有独立的列表。

解释:

# 错误示例:使用共享的可变默认值  
@dataclass  
class BadConfig:  
    station_codes: list[str] = []  

a = BadConfig()  
b = BadConfig()  

a.station_codes.append("ALERT")  
print(a_station_codes)  # ['ALERT']  
print(bstation_codes)  # ['ALERT']  // 这两个列表是相同的!!

最后是关于Kaggle设置的配置单例:

@dataclass
class KaggleConfig:
    """用于配置Kaggle数据集的发布相关设置。"""
    dataset_slug: str = "grimespoint/paris-flood-dataset"
    input_csv: str = "kaggle/input/datasets/{slug}/paris_flood_dataset.csv"
    output_dir: Path = field(default_factory=lambda: Path("kaggle/working/kaggle_dataset"))
    mock_output_dir: Path = field(default_factory=lambda: Path("mock_output"))
    output_filename: str = "paris_flood_dataset.csv"
    mock_output_filename: str = "mock_flood_dataset.csv"
    metadata_filename: str = "dataset-metadata.json"

    # 元数据信息
    title: str = "巴黎洪水数据集"
    keywords: list = field(default_factory=lambda: [
        "表格数据", "天气与气候", "环境研究", "欧洲", "时间序列分析"
    ])
    geospatial_coverage: str = "法国巴黎"
    update_frequency: str = "每周"
    license_name: str = "CC0-1.0"

    # 计算得出的字段(在__post_init__方法中生成)
    output_csv_path: Path = field(init=False)
    metadata_path: Path = field(init=False)

    def __post_init__(self):
        """在初始化之后计算这些派生字段的值。"""
        self.input_csv = self.input_csv.format(slug=self.dataset_slug)
        self.output_csv_path = self.output_dir / self.output_filename
        self.metadata_path = self.output_dir / selfmetadata_filename
        self.mock_output_filename = selfmock_output_dir / self.mock_output_filename

# 初始化各种配置对象——这些都是模块级别的单例对象
API_CONFIG = APIConfig()
STATION_CONFIG = StationConfig()
KAGGLE_CONFIG = KaggleConfig()

这种__post_init__的使用方式非常常见。output_csv_pathmetadata_path被标记为field(init=False),因此你不能通过构造函数直接设置它们的值。相反,这些字段的值会在__post_init__方法中根据其他字段的值来计算得出。

对于那些需要通过计算才能得到的值,应该采用这种处理方式:只在某个地方进行一次计算,这样在代码的其他部分每次需要使用这些路径时,就无需再次进行计算了。

有关这一设计模式的更多信息,可以参考Real Python关于数据类的指南,或者O’Reilly出版的《Fluent Python》一书中的数据类相关章节。

提示:自动生成的__repr__方法可以让你轻松地查看对象的所有字段及其值。只需调用print(API_config),就可以看到所有的信息,而无需编写任何格式化代码。在调试数据处理流程时,这个功能非常实用。

print(APIConfig)
# 输出结果如下:APIConfig/useMock=True, base_url='https://hubeau.eaufrance.fr/api/v2/hydrometrie/obs_elab', ...)

第四部分:数据提取步骤

在数据提取阶段,系统会从源系统中获取数据并将其读入内存。此时,你实际上有两个数据来源可以从中提取信息:一是之前已经发布过的现有数据集,二是通过Hub’Eau API获取的最新数据。

优雅的文件加载机制

设计逻辑

当这个处理流程首次运行时(此时还不存在任何数据集),会发生什么呢?如果采用简单的实现方式,程序很可能会因为遇到FileNotFoundError而崩溃。

其实解决方法很简单:load_csv()函数遵循了空对象模式。它不会抛出错误,而是返回一个空的DataFrame。这样一来,后续的所有处理函数都可以将“没有数据”和“有部分数据存在”的情况视为相同的情况,而无需进行任何特殊处理。

代码级详解

这个函数用于将CSV文件加载到DataFrame中。如果文件不存在,它会返回一个空的DataFrame,而不会导致程序崩溃。

low_memory=False这个参数告诉pandas要仔细地读取文件内容,从而避免出现类型判断错误。parse_dates=True则会让pandas尝试自动将包含日期信息的列转换为日期格式。delimiter=","表示文件是用逗号分隔的。

def load_csv(path: str) -> pd.DataFrame:
    """加载CSV文件;如果文件不存在,则返回一个空的DataFrame。

    参数:
        path (str): CSV文件的完整路径。

    返回值:
        pd.DataFrame: 如果文件存在,其中包含加载后的数据;如果文件不存在,返回一个空的DataFrame。

    异常情况:
        pd.errors ParserError: 如果CSV文件格式不正确,会抛出这个异常。
    """
    if os.path.exists(path):
        return pd.read_csv(path, low_memory=False, parse_dates=True, delimiter=",")
    return pd.DataFrame()   # 采用空对象模式,确保返回值的类型始终一致

为了测试这个功能,可以尝试使用一个不存在的文件路径来调用这个函数:

df_missing = load_csv("/tmp/does_not_exist.csv")
print(dfmissing.empty)  # True:程序没有崩溃

# 调用者也可以直接使用这种方式来判断结果是否为空:
if df_missing.empty:
    print("没有数据存在。将从头开始重新获取所有数据。")

最佳实践:函数中的每一条代码路径都应该返回相同类型的结果。如果一个函数有时返回DataFrame,有时返回None,那么调用者就必须在使用结果之前先检查它是否为None

无论数据是否存在,都直接返回DataFrame,这样会让代码更加简洁。在使用数据之前,你根本不需要去思考“我得到的是真实的数据吗?还是只是None?”这样的问题。

增量更新逻辑

设计逻辑

这就是第一部分中提到的“增量加载”机制的实际应用。在开始获取数据之前,首先要先问自己:“我目前已经拥有的最新数据是什么?我真的还需要再获取更多数据吗?”

策略:首先检查现有的数据是否已经包含了昨天的数据。

  • 如果已经包含:就直接跳过这次更新操作,因为数据集已经是最新的了。

  • 如果没有包含:就从已知数据的次日开始获取新数据。

现有数据:1月1日 – 1月15日  
昨天:1月19日  

决定:从1月16日开始获取数据(而不是从1月1日开始)。

为什么选择昨天而不是今天呢?因为源系统可能还没有完成今天的数据更新。Hub'Eau提供的“处理后的观测数据”其实是每天汇总一次后的结果,所以最稳妥的比对方式就是使用最近一个已完成更新的数据。

代码级讲解

determine_update_range()函数会检查最新的数据保存日期。它会告诉你“你的数据已经是最新的了”,或者“应该从这个日期开始下载新数据”。

  1. 该函数首先会使用现有的数据

    • existing是一个已经包含了一些数据的pandas DataFrame对象。
  2. 如果根本没有数据

    • if existing.empty:

    • 如果这个DataFrame中的行数为0,那么函数会返回:“还没有任何数据被保存下来,因此需要下载所有数据。”

    • 函数会返回以下结果:

      • True = 需要更新数据

      • STATION_CONFIG.earliest_date = 应从最早允许的日期开始下载数据

    • 找到包含日期的列

      • 代码会检查哪一列包含了日期信息:

        • 首先会尝试使用"date_obs_elab"这个列名

        • 如果找不到,就会尝试使用"record_date"

      • 如果这两列都不存在,函数就会抛出错误,因为它不知道应该使用哪一列来获取日期信息。

    • 将包含日期的列转换为真正的日期对象

      • pd.to_datetime(...)这个函数可以将该列转换成pandas能够处理的日期对象。

      • errors="coerce"这个参数表示:如果遇到无效的日期值,程序会将其设置为缺失值,而不会出现错误。

    • 找出数据中的最新日期

      • last_day = s.max().date()

      • 这个操作会获取数据集中最新的日期值。

    • 将这个最新日期与昨天进行比较

      • yesterday = date.today() - timedelta(days=1)

      • 函数会检查数据中是否已经包含了昨天的日期。

    • 如果数据已经是最新的

      • 如果last_day >= yesterday,函数会返回:

        • False = 不需要更新数据

        • None = 不需要指定开始下载的日期

    • 如果数据还落后于最新日期

      • 函数会将next_day设置为最后一个保存日期的次日。

      • 然后函数会返回:

        • True = 需要更新数据

        • 并且会以字符串形式返回下一个需要下载的日期,例如"2026-07-12"

def determine_update_range(existing: pd.DataFrame) -> Tuple[bool, Optional[str]]:
    """判断是否需要更新数据,以及应该从哪一天开始更新。

    处理逻辑:
    1. 检查现有数据是否包含昨天的日期;
    2. 如果包含,则不需要更新;
    3. 如果不包含,则从前一天开始获取数据。

    返回值:
        Tuple[是否需要更新, 开始更新的日期]
    """
    # 情况1:磁盘上还没有任何数据
    if existing.empty:
        print("没有找到现有数据。将从最早的时间点开始获取所有数据.")
        return True, STATION_CONFIG.earliest_date

    # 如果现有数据的列中包含“date_obs_elab”或“record_date”,则使用相应的列名;
    # 否则,会抛出KeyError异常。
    if "date_obs_elab" in existing.columns:
        record_colname = "date_obs_elab"
    elif "record_date" in existing.columns:
        record_colname = "record_date"
    else:
        raise KeyError("缺少日期列:应该包含‘date_obs_elab’或‘record_date’")

    # 将现有数据中的日期转换为datetime类型
    s = pd.to_datetime(existing[record_colname], errors="coerce")
    # 获取最新的日期
    last_day = s.max().date()
    # 计算昨天的日期
    yesterday = date.today() - timedelta(days=1)

    # 情况2:数据已经包含了昨天的日期或更晚的日期
    if last_day >= yesterday:
        print("\n数据集已经包含了昨天的日期或更晚的日期,因此不需要更新.")
        return False, None

    # 情况3:需要获取缺失的数据段
    next_day = (last_day + pd.Timedelta(days=1))
    print(f"\n将从{next_day}开始获取数据.")
    return True, next_day.isoformat()

有几点值得注意。

首先,该函数会检查两种可能的列名:原始API名称`date_obs_elab`,或者已经更名的英文名称`record_date`。它并不只假设存在其中一种列名。因此,无论你是使用刚获取的原始数据,还是从磁盘加载的已处理CSV文件来调用这个函数,它都能正常工作。

errors="coerce"这一设置在这里也会被用到,我们将在第5部分进一步探讨这一点。任何pandas无法解析的日期都会被标记为`NaT`(表示“非时间值”),而不会引发异常。

该函数的返回类型是`Tuple[bool, Optional[str]]`。这种元组包含了两个相关的结果:是否应该更新数据?以及应该从何时开始更新?这种方式比直接返回两个独立的值,或者更糟糕的是返回一个根据上下文含义不同的模糊值,要好得多。

最佳实践:使用`Tuple`作为返回类型(如果需要处理更多字段,也可以使用小型数据类或`NamedTuple`),将相关的结果整合在一起。同时要清楚地说明每个元素的含义。如果你在不同代码路径中返回了不同类型的值,却没有相应的文档说明,那么就很容易引发混淆和错误,比如“为什么有时候这个值的类型是`None`,而有时候又是字符串?”

我们可以通过以下三种场景来快速验证这个函数的正确性:

# 测试1:没有现有数据
should_update, start_date = determine_update_range(pd.DataFrame())
# → True, "1900-01-01"

# 测试2:存在部分旧数据(仅覆盖1月10日至15日的记录)
# → True, "2026-01-16" (最后一个已知日期的次日)

# 测试3:最近的数据已经包含了昨天的信息
# → False, None

模拟API接口

设计逻辑

在实际生产环境中,获取数据意味着要发起真实的HTTP请求:

requests.get(
    "https://hubeau.eaufrance.fr/api/v2/hydrometrie/obs_elab",
    params={"code_entite": "F700000109", "size": 20000, ...}
)

真实的API调用会带来一些实际挑战。不过,在本教程中,我们首先会使用模拟接口生成器来构建和测试所有逻辑。这个生成器返回的数据格式与真实API完全一致。只有当所有的逻辑都通过测试后,才会切换到真实的API端点进行测试。

这种技术在这个项目之外的场景中也非常有用。在开始实际开发之前,先使用模拟数据或固定测试用例来验证你的数据处理逻辑,这样就可以避免同时面对“我的解析是否出错?”以及“当前网络连接是否不稳定?”这类问题。

代码层级详解

generatemock_api_data()这个函数会为某个监测站点生成虚假的样本数据,每天生成一条记录,数据起始时间为`start_date`。

其工作原理如下:

  • 它会将`start_date`转换为实际日期。

  • 然后它会循环`num_days`天。

  • 对于每一天,都会生成一个模拟的观测数据字典。

  • 会在水位数值上添加一些随机变化,使数据看起来更真实。

  • 会随机选择状态、水质等级以及测量方法的相关标签。

  • 最后返回这些字典组成的列表。

def generateMockApiData(station_code: str, start_date: str, num_days: int = 10) -> List[Dict]:
    """生成用于演示的、逼真的模拟API数据。

    该函数模拟HubEau API会返回的数据格式:即一系列观测数据字典。"""
    start = pd.to_datetime(start_date).date()
    records = []

    validation statuses = ["已验证数据", "原始数据", "预验证数据"]
    qualities = ["良好", "未达标", "可疑"]
    methods = ["实测数据", "计算数据", "专家评估数据"]

    for i in range(num_days):
        obs_date = start + timedelta(days=i)
        base_level = 5500 + int(station_code[-2:])  # 不同站点这个数值会有所不同
        noise = random.randint(-200, 200)
        water_level = base_level + noise

        record = {
            "code_site": "mock_" + station_code[1:],
            "code_station": "mock_" + station_code,
            "date_obs_elab": obs_date.isoformat(),
            "resultat_obs_elab": water_level,
            "date_prod": (obs_date + timedelta(days=1)).isoformat(),
            "code_statut": "1",
            "libelle_statut": random.choice(validation statuses),
            "code_methode": "1",
            "libelle_methode": random.choice(methods),
            "code_qualification": "1",
            "libelleQUALIFICATION": random.choice(qualities),
            "longitude": 2.3522 + random.uniform(-0.01, 0.01),
            "latitude": 48.8566 + random.uniform(-0.01, 0.01),
            "grandeur_hydro_elab": "mock_HIXnJ",
        }
        records.append(record)

    return records

实际数据获取流程

设计逻辑

有两个函数负责数据的提取工作。根据单一职责原则,这两个函数被刻意分开了:每个函数都应该只负责一个特定的功能。

fetch_all_data() (协调器)
    ├── fetch_single_station_data(station_1) ← 负责处理所有复杂逻辑
    ├── fetch_single_station_data.station_2) ← 负责处理所有复杂逻辑
    └── fetch_single_station_data(station_n) ← 负责处理所有复杂逻辑
  • fetch_single_station_data()函数负责处理与每个站点相关的一切复杂细节:分页、游标控制、停止条件以及网络错误处理等。

  • fetch_all_data()函数则不涉及这些细节。它只是遍历所有站点,然后将数据提取任务委托给相应的函数进行处理。

这种分工方式有两个好处。首先,如果你想为某个站点修改分页策略,完全不需要修改协调器的代码;其次,如果你想并行执行数据获取操作(例如使用concurrent.futuresasyncio),协调器就是你唯一需要修改的部分。

代码层级解析

fetch_single_station_data()函数会从模拟测试数据或真实的API中逐页获取站点数据,直到收集到所有所需的信息为止。

  • 如果 useMock=True,那么它会使用虚拟数据来代替调用真实的API。

    1. 它会调用generate_mock_api_data()函数

    2. 然后将得到的结果转换成DataFrame格式

    3. 会将日期字段转换为pandas支持的日期类型

    4. 最后返回这个DataFrame

  • 如果 useMock=False,那么它就会进行真实的API请求

    1. 会创建一个可重复使用的HTTP会话

    2. start_date开始获取数据

    3. 会反复向API请求数据页面

    4. 在以下情况下停止请求:

      • API没有返回任何数据

      • 最新的日期已经到达昨天

      • 获取到的数据页数少于预期

      • 发生了网络错误

    5. 会将所有获取到的数据合并成一个DataFrame

    6. 如果没有获取到任何数据,就会返回一个空的DataFrame

def fetch_single_station_data(station_code: str, start_date: str, useMock: bool = True) -> pd.DataFrame:
    """从虚拟API或真实API端点获取某个站点的所有水文数据。

    在生产环境中(采用基于游标的分页机制):
    - 每次请求只获取max_per_page条记录
   > 会持续请求,直到没有新数据或到达昨天为止
   - 在以下情况下停止请求:没有数据返回 | 最后获取的日期大于等于昨天 | 获取到的数据页数不足
   - 能够优雅地处理网络错误
    """
    if useMock:
        data = generatemock_api_data(station_code, start_date, num_days=7)
        page_df = pd.DataFrame(data)
        page_df["date_obs_elab"] = pd.to_datetime(
            page_df["date_obs_elab"], errors="coerce").dt.normalize()
        return page_df

    # 真实的数据获取逻辑
    else:
        session = requests.Session()  # 在多次请求中重用TCP连接
        frames = []
        cursor = start_date

        while True:
            params = {
                "code_entite": station_code,
                "grandeur_hydro_elab": API_CONFIG(metric),
                "date_debut_obs_elab": cursor,
                "size": API_config.max_per_page,
            }

            try:
                response = session.get(
                    API_CONFIG.base_url,
                    params=params,
                    timeout=API Config.timeout_seconds  # 建议始终设置超时时间
                )
                response.raise_for_status()
            except requests.RequestException as e:  # 避免因为某个站点的请求失败而影响整个流程
                print(f"在获取站点{station_code}的数据时发生错误:{e}")
                break

            data = response.json().get("data", [])
            if not data:  # 如果没有数据返回,说明已经没有更多数据可获取了
                break

            page_df = pd.DataFrame(data)
            page_df["date_obs_elab"] = pd.to_datetime(
                page_df["date_obs_elab"], errors="coerce").dt.normalize()
            frames.append(page_df)

            last_page_date = page_df["date_obs_elab"].max()
            yesterday = date.today() - timedelta(days=1)

            # 防止无限循环
            if pd.isna(last_page_date) or last_page_date.date() >= yesterday:
                break

            cursor = (last_page_date + pd.Timedelta(days=1)).strftime("%Y-%m-%d")

            if len(data) < API_CONFIG.max_per_page:
                break

    if frames:
        return pd.concat(frames, ignore_index=True)
    return pd.DataFrame()

分页循环实际上是如何工作的呢?

让我们一步步来分析这个过程。这是整个脚本中逻辑最为复杂的部分。

  1. 发送一个请求,其中包含参数 date_debut_obs_elab=cursor,即“从这一天开始获取记录”。

  2. 如果请求失败(出现 requests.RequestException 异常),就记录错误信息并立即跳出循环。这样就可以避免因为某个站点的网络问题而影响整个数据采集流程。

  3. 如果响应中根本没有数据,说明你已经获取到了所有需要的记录:此时也跳出循环

  4. 否则,记录下当前页面上显示的最新日期(即 last_page_date)。

  5. 如果这个最新日期已经大于或等于昨天,说明你已经获取到了所有需要的记录:此时也跳出循环

  6. 否则,将游标移动到 last_page_date + 1天 的位置,然后重新开始循环以获取下一页的数据。

  7. 作为安全措施:如果页面返回的记录数量少于 max_per_page,这也说明你已经到达了数据的末尾。因为API只有在数据不足时才会返回不完整的页面数据,所以此时也跳出循环

最后这个检查步骤(第7步)其实是一种经典的分页终止策略。你并不总是需要从API那里获取next_page令牌;如果一个页面应该包含20,000条记录,但实际只返回了4,213条记录,那就说明已经没有更多的数据可以获取了。

这里有三个经过慎重考虑的特定设计决策:

session = requests.Session()

Session对象可以在对同一主机发送多次请求时重用底层的TCP连接。这样就可以避免每次请求都需要重新建立TCP/TLS连接,从而提高效率,同时也能减少对服务器的压力。

最佳实践:每当在循环中多次调用requests.get()来访问同一个主机时,都应该使用Session对象。

except requests.RequestException as e:

RequestExceptionrequests库中所有异常类型的基类。无论是超时错误、连接错误,还是由raise_for_status()引发的HTTP错误,都可以通过捕获这个基类来统一处理。这样,无论遇到什么样的网络问题,都能以相同的方式进行处理:记录错误信息,停止获取该站点的数据,然后继续执行后续操作。

timeout=API_CONFIG.timeout_seconds

默认情况下,requests库是从不设置超时时间的。如果没有明确指定超时时间,那么某个服务器如果出现故障,就可能会无限期地阻塞整个数据采集流程。requests的高级使用指南建议对外部服务器的请求必须设置超时时间。

最佳实践: 将所有的外部I/O操作都包裹在`try/except`语句中。无论发生什么情况,都要确保程序能够优雅地处理错误并继续运行。记录下错误信息,让整个流程能够恢复正常或继续执行后续步骤。这样就可以避免因为某个出问题的请求而导致整个定期执行的作业失败。

现在,这个负责协调数据获取的函数`fetch_all_data()`的设计已经变得简单多了:

  1. 遍历所有站点代码。

    • 获取每个站点的数据。

    • 保留那些非空的结果。

  2. 将这些数据合并成一个大的DataFrame。

def fetch_all_data(start_date: str, useMock: bool = True) -> pd.DataFrame: """这个函数负责获取所有配置好的站点的数据。""" frames = [] for station_code in STATION_CONFIG.station_codes: print(f"正在获取站点{station_code}的数据...") df_station = fetch_single_station_data(station_code, start_date, useMock=useMock) if not df_station.empty: print(f>获取到了{len(df_station)}条记录") frames.append(df_station) else: print(f>没有找到数据…) if frames: return pd.concat(frames, ignore_index=True) return pd.DataFrame()

就是这么简单。只需要一个循环和`pd.concat`函数而已。所有那些复杂的逻辑都被隐藏在较低层次的代码中,正好符合SRP的设计原则。

第5部分:数据转换环节

在这个阶段,原始的、刚刚获取到的数据会被转化为可以直接发布的形式。整个数据转换过程的设计理念是:使用许多小型且功能单一的函数来完成具体的处理任务。每个函数都只接收一个DataFrame作为输入,然后返回一个新的DataFrame作为结果,而且不会产生任何副作用。这样就可以避免使用那些功能过于复杂的“万能函数”。

(提取数据) 原始API数据 ↓ (转换数据) 1. 类型解析(将日期、数值等转换为适当类型) 2. 列名修改(将法语列名改为英语) 3. 对分类变量进行重新映射 4> 计算衍生列(例如洪水警报标志) 5> 重新排序列目 6> 对数据进行排序并重置索引 ↓ (加载数据) 可以直接发布的数据集

为什么要把这个过程分成六个步骤,而不是使用一个庞大的函数呢?因为每个步骤都是独立可测试的,也可以单独进行修改。如果在凌晨3点某个定时任务执行时出现了问题,那么可以单独调试每一个函数,从而准确地找出是哪个环节出了故障。这样就不需要去分析那些长达200行的复杂函数了。

类型解析与优雅的强制转换

设计逻辑

从CSV文件或JSON API中获取到的数据,在默认情况下都是以字符串的形式存在的。Pandas需要这些数据具有明确的类型,才能对日期进行时间顺序排序、进行数值运算(比如`last_date + timedelta(days=1)`),或者比较不同的数值(比如`water_level > 6000`)。如果类型不匹配,就会导致程序出现错误,从而影响整个数据处理流程。例如,如果有一行数据的内容是`"N/A"`,或者日期被截断了,又或者数值对齐不正确,那么严格的类型解析器就会抛出异常,从而导致整个任务失败。

最佳实践: 下面这两个转换函数都使用了 `errors="coerce"` 选项,这样当 pandas 无法解析某些值时,这些值会变成 `NaT`(非时间类型)或 `NaN`(非数字类型),而不会引发错误。这种处理方式被称为“优雅强制转换”。

代码级详解

convert_to_date()convert_to_numeric() 是一些辅助函数,它们的作用是确保某些列的数据类型符合要求。

  • 这两个函数都是先通过 `df.copy()` 创建一个副本,这样就不会修改原始的 DataFrame。

  • errors="coerce" 这个选项意味着无效的数据会变成缺失值,而不会导致程序崩溃。

对于 convert_to_date() 来说:

  1. 首先通过 `df.copy()` 创建一个副本,这样原始表格就不会被修改。

  2. 然后循环遍历 `columns` 中列的名字。

    • if col in df.columns 这一行代码会检查该列是否确实存在,然后再尝试进行转换。

    • pd.to_datetime(...) 会将像 “2026-07-12” 这样的文本转换为真正的 pandas 日期/时间类型。

    • errors="coerce" 使得无效的数据会变成 `NaT`,而不会引发错误。

    • .dtnormalize() 会去除时间部分,只保留午夜时的日期值。

对于 convert_to_numeric() 来说:

  1. 同样也是先通过 `df.copy()` 创建一个副本。

  2. 然后循环遍历 `columns` 中列的名字。

    • if col in df.columns 会检查该列是否确实存在,然后再尝试进行转换。

    • pd.to_numeric() 会将像 “12.5” 这样的文本转换为数字类型。

    • errors="coerce" 使得无效的数据会变成 `NaN`,而不会导致程序崩溃。

def convert_to_date(df: pd.DataFrame, columns: List[str]) -> pd.DataFrame:
    """将指定的列转换为 pandas 的日期时间类型。"""
    df = df.copy()  # 绝不要修改原始数据!
    for col in columns:
        if col in df.columns:
            df[col] = pd.to_datetime(df[col], errors="coerce").dtnormalize()
    return df


def convert_to_numeric(df: pd.DataFrame, columns: List[str]) -> pd.DataFrame:
    """将指定的列转换为数值类型。"""
    df = df.copy()
    for col in columns:
        if col in df.columns:
            df[col] = pd.to_numeric(df[col], errors="coerce")
    return df

为什么这样做有用呢?

  • 这样就可以对数据进行处理了:比如对其进行分组、汇总等等。

  • 日期类型的数据现在可以像普通日期一样进行排序和筛选了。

  • 数字类型的数据在计算平均值、总和或进行比较时也能正常使用。

  • 这种处理方式还能避免因为数据类型混杂而导致的错误,比如 “12” 和 “12.0” 这种情况。

pd.to_datetimepd.to_numeric都是pandas提供的官方函数,它们都带有errors参数。默认情况下,该参数被设置为"raise",这意味着遇到无效输入时程序会抛出异常。其他可选值包括"coerce"(将无效数据替换为null)和"ignore"(不对其进行任何处理)。如需查看完整的参数列表,请参阅pandas的to_datetime函数文档以及to_numeric函数文档

让我们用一个包含无效数据的示例来测试这些函数的功能:

messy_df = pd.DataFrame({
    "date_obs_elab": ["2026-01-15", "2026-01-16", "not a date", None],
    "resultat_obs_elab": [5800.0, "5900", "N/A", None],
})

type_safe_df = convert_to_date(messy_df, ["date_obs_elab"])
type_safe_df = convert_to_numeric(type_safe_df, ["resultat_obs_elab"])

# "not a date"  → NaT
# "N/A"         → NaN
# 处理流程能够正常进行,没有任何错误发生。

最佳实践:不要让其中一行无效数据导致整个处理流程失败。应将无效数据替换为null,而不是抛出异常。如果以后需要检查数据质量,可以单独标记或记录这些null值。

这种处理方式其实是一种有意识的权衡:在流程的连续性数据的严格验证之间做出选择。对于那些需要自动执行且无人监控的处理任务来说,这种做法通常是正确的。

df.copy():基于惯例实现数据不可变性

再看看上面这两个函数的开头部分:df = df.copy()。这条代码会出现在流程中每一个转换函数的开始处,这绝非偶然。

Python中的DataFrame属于可变对象,它们是通过引用传递的。如果某个函数直接修改了df的内容而没有先创建一个副本,那么调用者的原始DataFrame也会被改变。这就是典型的副作用,这种错误很容易导致程序出现混乱。因此,最好先调用.copy(),这样每个函数处理的输出都会是全新的对象,调用者提供的原始数据也会保持不变

最佳实践:应将DataFrame视为不可变输入。修改数据时应该创建一个新的DataFrame,而不是直接修改原对象。即使这样做会消耗少量的内存或CPU资源,但对于那些规模不是很大的处理流程来说,这种做法带来的调试便利性几乎总是值得的。

最后,还有一个自动检测功能的便捷函数,它可以结合这两个底层函数来使用。这个函数会扫描文本列,判断其中是否包含日期或数字,然后自动调用相应的解析函数进行处理。

auto_convert_columns()能够自动识别哪些列是日期或数字,并据此自动调整这些数据的类型。

它的具体工作原理是……

整个过程从两个空列表开始:

  • datetime_cols 用于存储日期类型的列

  • numeric_cols 用于存储数值类型的列

程序会遍历 DataFrame 中的每一列。如果某列已经是日期或数值类型,就会直接跳过它;如果该列是文本类型(如 objectstring),则会查看该列前 10 个非空样本值。

首先,程序会尝试将这些样本值解读为日期:

  • 如果解析成功,该列就会被添加到 datetime_cols

如果解析失败,就会尝试将它们解读为数值:

  • 如果解析成功,该列就会被添加到 numeric_cols

最后,程序会使用 convert_to_date() 将所有日期类型的列转换为正确的格式,然后使用 convert_to_numeric() 将所有数值类型的列转换为正确的格式。

def auto_convert_columns(df: pd.DataFrame) -> pd.DataFrame:
    """自动检测并将日期类型及数值类型的列转换为正确的格式。"""
    datetime_cols = []
    numeric_cols = []

    for col in df.columns:
        if pd.api.types.is_datetime64_any dtype(df[col]):
            continue
        if pd.api_types.is_numeric_dtype(df[col]):
            continue

        if pd.apitypes.is_objectdtype(df[col]) or pd.api_types.is_stringatype(df[col]):
            sample = df[col].dropna().head(10)
            if len(sample) == 0:
                continue

            try:
                pd.to_datetime(sample, errors='raise', format='mixed')
                datetime_cols.append(col)
                continue
            except (ValueError, TypeError):
                pass

            try:
                pd.to_numeric(sample, errors='raise')
                numeric_cols.append(col)
                continue
            except (ValueError, TypeError):
                pass

    df = convert_to_date(df, datetimecols)
    df = convert_to_numeric(df, numeric_cols)
    return df

需要注意的是,这里的内部 try/except 块使用了 errors='raise' 这一设置。这与之前的强制转换策略正好相反,因为这里只针对每一列的前 10 个样本值进行测试。

这其实只是一个类型检测步骤,并非最终的转换操作。它只是通过少量的样本数据来判断“这一列属于日期类型、数值类型,还是两者都不是”。之后,才会对整个列进行正式的转换操作,这些转换任务分别由 convert_to_date()convert_to_numeric() 完成。两种不同的错误处理策略,对应着两种不同的处理流程。

具有双向映射关系的数据结构转换

设计逻辑

Hub'Eau API 返回的是法语列名和法语分类值,例如 code_station"Donnée validée"。针对国际用户设计的数据集应该使用英语进行表述。如果随意在代码中重新命名列名,那么法语和英语之间的映射关系就会分散在代码的各个部分,到时候如果要反向转换这些名称,就会遇到麻烦。

解决方法:在模块的开头定义一个权威的映射关系,然后根据这个映射关系推导出其他所有内容。

代码级讲解

API_TO_EN将API字段名称映射为英文对应名称。

# 主要映射关系:法文API字段名与英文字段名的对应关系
API_TO_EN = {
    "code_site": "location_code",
    "code_station": "station_code",
    "date_obs_elab": "record_date",
    "resultat_obs_elab": "water_level_mm",
    "date_prod": "data_production_date",
    "code_statut": "validation_status_code",
    "libelle_statut": "validation_status",
    "code_methode": "production_method_code",
    "libelle_methode": "production_method",
    "code_qualification": "quality_code",
    "libelleQUALIFICATION": "quality_assessment",
    "longitude": "longitude",
    "latitude": "latitude",
    "grandeur_hydro_elab": "hubeau_elab_code",
}

# 反向映射关系:英文名称对应法文名称(自动生成)
EN_TO_API = {v: k for k, v in API_TO_EN.items()}

第二行代码EN_TO_API其实只是一个快捷方式:它只是将API_TO_EN中所有的键值对进行互换。这种映射关系是通过字典推导式生成的。关键的设计点在于:EN_TO_API并不是人工维护的,而是通过其他代码自动生成的。如果你在API_TO_EN中添加、删除或修改任何条目,那么每次模块运行时,EN_TO_API都会自动更新。

在整个代码库中,只有一个地方需要对数据结构进行修改。

对于分类型的(而不仅仅是字段名称),也会采用相同的处理方式:

CATEGORICAL_MAPPINGS = {
    "validation_status": {
        "Donnée validée": "validated",
        "Donnée brute": "raw",
        "Donnée pré-validée": "pre-validated",
    },
    "quality_assessment": {
        "Bonne": "good",
        "Non qualifiée": "unqualified",
        "Douteuse": "dubious",
    },
    "production_method": {
        "Calculée": "calculated",
        "Mesurée": "measured",
        "Expertisée": "expert-reviewed",
    },
}

下面的函数会依次应用这些映射关系,从而标准化字段名称和分类值。

rename_to_english():

  1. 如果DataFrame为空,它会立即返回一个副本。

  2. 它会生成一个列表,列出那些需要从API名称更改为英文名称的字段。

    • 只有当满足以下条件时,才会对某个字段进行重命名:

      • 该字段的旧名称确实存在,

      • 并且新的名称还没有被使用过。

    • 完成重命名后,它会返回修改后的DataFrame。

rename_to_api_schema():

  1. 对于空数据,也会返回一个副本。

  2. 它的作用是将英文名称转换回API名称。

  3. 最后会返回修改后的DataFrame。

(如果你需要以API原始格式将数据传回,这个方法会非常有用。)

apply_categoricalMappings():

  1. 它会创建一个副本,因此原始的DataFrame不会被修改。

  2. 对于CATEGORICAL_MAPPINGS中的每一列,它都会使用对应的映射关系来替换其中的值。

    • 例如:"Donnée validée"会被替换成"validated"

    • .fillna(df[col_name])这个方法会在映射中找不到某个值时,保留该值的原始值。

  3. 最后返回修改后的DataFrame。

def rename_to_english(df: pd.DataFrame) -> pd.DataFrame:
    """将API列名重命名为英文格式的名称。"""
    if df.empty:
        return df.copy()

    columns_to_rename = {}
    for src, dst in API_TO_EN.items():
        if src in df.columns and dst not in df.columns:
            columns_torename[src] = dst

    return df.rename(columns=columns_to Rename)

def rename_to_api_schema(df: pd.DataFrame) -> pd.DataFrame:
    """将英文列名重新恢复为API格式的名称(即执行逆向操作)。"""
    if df.empty:
        return df.copy()

    columns_to_rename = {k: v for k, v in EN_TO_API.items() if k in df.columns}
    return df.rename(columns=columns_to Rename)
def apply_categorical_mappings(df: pd.DataFrame) -> pd.DataFrame:
    df = df.copy()
    for col_name, mapping in CATEGORICAL_MAPPINGS.items():
        if col_name in df.columns:
            df[col_name] = df[col_name].map(mapping).fillna(df[col_name])
    return df

这里有两条值得注意的注意事项:

对部分数据的处理能力:这两个重命名函数只会修改输入数据中确实存在的列(即if src in df.columns这个条件成立的情况)。

如果一个重命名函数假设所有被映射的列都一定存在于输入数据中,那么当它在数据结构不完整或格式不同的DataFrame上被调用时,就会导致程序崩溃。

.map(mapping).fillna(df[col_name])这个方法会用映射字典中的值来替换数据中的各个元素;而对于那些在字典中找不到的元素,它会将其设置为NaN。而.fillna(df[col_name])这个操作则会在映射关系不适用的地方恢复元素的原始值。这样一来,那些意外的分类数值就能保持原有的状态,而不会被悄悄地设置为null。这是一个细微但非常重要的设计选择,能够确保程序的稳定性。

关于这种.map()后再使用.fillna()的处理方式,可以在Stack Overflow上的这个链接中找到详细讨论。

一个简单的测试就可以证明这种双向映射功能确实是有效的:

sample_api_df = pd.DataFrame({"code_station": ["F700000109"], ...})
renamed_df = rename_to_english(sample_api_df)          # 法文列名 → 英文列名
reversed_df = rename_to_api_schema(renamed_df)          # 英文列名 → 法文列名
# reversed_df.columns.tolist() == sample_api_df.columns.tolist()  → True

计算洪水警报

add_derived_columns()只是一个简单的函数。它的作用就是计算flood_alert这个值。如果水位超过了在Config中设置的洪水阈值,那么结果就会是True;否则,结果就是False

def add_derived_columns(df: pd.DataFrame) -> pd.DataFrame:
    """根据原始数据添加计算出的列。"""
    df = df.copy()

    if "water_level_mm" in df.columns:
        df["flood_alert"] = df["water_level.mm"] > STATION_CONFIG.flood_threshold_mm

    return df

列排序

这里有一个虽然很简单,但对用户来说却非常重要的细节:一旦确定了首选的列排序顺序,就可以在所有地方重复使用这个顺序。COLUMN_ORDER变量就保存了这种优选的列序列。

order_columns()函数会重新排列DataFrame中的列,使得这些优先列排在前面,而其他多余的列则会被放在后面。

COLUMN_ORDER = [
    # 主要标识符和测量数据
    "station_code", "record_date", "water_level_mm", "flood_alert",
    # 与观测数据相关的元信息
    "hubeau_elab_code", "data_production_date",
    "validation_status_code", "validation_status",
    "production_method_code", "production_method",
    "quality_code", "quality_assessment",
    # 地理位置信息(相对次要)
    "location_code", "longitude", "latitude",
]

def order_columns(df: pd.DataFrame) -> pd.DataFrame:
    """按照首选顺序重新排列列。"""
    present_cols = [c for c in COLUMN_ORDER if c in df.columns]
    othercols = [c for c in df.columns if c not in present_cols]
    return df[present_cols + other cols]

othercols这个机制起到了“安全网”的作用:任何没有在COLUMN_ORDER中明确列出的列,都会被保留在数据集的末尾,而不会被自动删除。

列排序应遵循以下优先规则:

  • 将标识符和关键字段放在最前面。

  • 把最重要、使用频率最高且最稳定的列排在前面。

  • 将相关的字段放在一起,以便数据表看起来更清晰、更易于阅读。

  • 将可选或很少使用的字段放在数据集的末尾。

去重处理

设计逻辑

在数据处理过程中应随时去除重复的数据。在增量式数据处理流程中,不同的日期范围或记录可能会重叠:

  • 由于API返回的数据可能延迟到达,或者数据会在之后才被最终确定,因此同一天的事情可能会被多次获取。

  • 在同一分页响应的不同页面上,也可能会出现相同的观测数据。

  • 如果不进行去重处理,这些情况会导致数据集中逐渐积累重复的行。回想一下第一部分的内容,正是这种去重特性使得数据处理流程具有幂等性:无论运行一次还是五次,最终得到的数据集都是一样的。

    要解决这个问题,首先需要明确什么因素才能使一条记录变得唯一。因此,我们需要定义一个复合键:
    key = (station_code, observation_date, water_level_value)

    为什么选择这种特定的组合方式呢?从物理角度来看,一个传感器会报告某一天的日最高水位数值,而这个数值应该是唯一的。

    • 不同站点的记录显然不会彼此重复。

    • 不同日期的两次测量数据也不会重复。

    • 还有一个更微妙的问题:如果同一个站点在同一天报告了不同的数值,那么这仍然应该被视为一次独立的观测数据。例如,这个数值可能是经过修正后的结果,而不是可以被直接忽略的重复数据。

    代码级讲解

    create_dedup_key()通过将三个部分用下划线连接起来,为每一行生成一个唯一的文本标识符:

    • 站点代码

    • 日期,格式为YYYY-MM-DD

    • 水位数值

    示例标识符:

    "F700000109_2024-01-15_5800.0"
    def create_dedup_key(df: pd.DataFrame) -> pd.Series:
        """根据站点代码、日期和水位数值生成唯一的去重标识符。
    
        标识符格式:"station_code_YYYY-MM-DD_value"
        例如:"F700000109_2024-01-15_5800.0"
        """
        parts = []
    
        if "code_station" in df.columns:
            parts.append(df["code_station"].astype(str))
    
        if "date_obs_elab" in df.columns:
            parts.append(df["date_obs_elab"].dt.strftime("%Y-%m-%d"))
    
        if "resultat_obs_elab" in df.columns:
            parts.append(df["resultat_obs_elab"].astype(str))
    
        if not parts:
            return pd.Series(index=df.index, dtype="object")
    
        return pd.Series(
            ["_".join(row) for row in zip(*parts)],
            index=df.index
        )
    

    标识符的生成过程如下: `parts`最终会成为一个列表,其中每个元素都是一个Series对象,分别对应站代码、日期和水位数值这三个组成部分,而且每个Series的对象长度都与原始DataFrame相同。

    让我们从内到外仔细理解最后这一行代码:

        return pd.Series(["_".join(row) for row in zip(*parts)], index=df.index)

    zip(*parts)会将这个列列表转换成按行排列的元组序列。例如,它会生成`(station_1, date_1, value_1)`、`(station_2, date_2, value_2)`等等。然后,列表推导式会将这些元组用下划线连接起来,从而为每一行生成一个唯一的字符串标识符。

    这种“列列表 → zip操作 → 按行排列的元组序列”的处理方式,是一种常见且高效的方法,可以用来将多个Series合并成一个新的Series对象,而无需使用`.apply(lambda row: ..., axis=1)`这样的代码。相比之下,在pandas中,按行进行的`.apply`操作通常会运行速度较慢,而向量化字符串操作则要高效得多。

    实际的去重操作是在remove_duplicates()这个函数中完成的。它会从"new"(新获取的数据)中删除那些已经存在于"existing"(历史数据)中的行。

    具体步骤如下:

    1. 如果其中一个数据表为空,那么它就会直接返回new

    2. 首先会使用auto_convert_columns()来确保两个DataFrame中的数据类型一致,这样日期和数字才能被正确地进行比较。

    3. 然后会为两个表格中的每一行生成一个去重键,这个过程是通过create_dedup_key()完成的。

    4. 接下来会检查新获取的数据中哪些键在历史数据集中并不存在,从而只保留真正新的数据行。

    5. 最后返回过滤后的结果,并保留new中的原始列。

    def remove_duplicates(existing: pd.DataFrame, new: pd.DataFrame) -> pd.DataFrame:
        """从'new'中删除那些已经存在于'existing'中的行。"""
        # 如果任意一个数据表为空,就无需进行任何操作
        if existing.empty or new.empty:
            return new.copy()
    
        # 确保两个数据表中的数据类型一致,以便正确比较
        existing_std = auto_convert_columns(existing)
        new_std = auto_convert_columns(new)
    
        # 生成去重键
        existing_keys = set(create_dedup_keyexisting_std).dropna())
        new_keys = create_dedup_key(new_std)
    
        # 创建布尔掩码:当新数据行确实是新的时,该掩码值为True
        mask = ~new_keys.isin(existing_keys)
    
        # 通过原始的(未标准化的)'new'数据表来提取最终结果,从而保留所有列
        result = new.iloc[new_keys[mask].index].copy()
        return result
    

    有两条关于性能和稳定性的要点值得注意:

    首先,existing_keys实际上是一个集合,而不是一个列表。在Python中,对集合进行成员判断操作(使用in.isin())的平均时间复杂度为O(1),因为这种操作会利用快速的哈希查找机制来直接判断某个元素是否已经存在;而如果使用列表,时间复杂度则为O(n),因为需要逐个检查每个元素,而且当数据集规模较大时,效率会显著降低。

    对于那些包含数万行数据的数据集来说,每次运行处理流程时都会进行这样的判断操作,因此这种性能差异是非常重要的。这正是选择合适的数据结构来解决问题的一个典型例子。有关Python中列表集合在查找操作中的性能对比,可以参考这个讨论;而关于使用.isin()进行过滤操作的性能优化方法,也可以参考这篇指南

    简单的规则是:当你需要快速判断某个元素是否存在于数据集中时,应该使用集合;而当你需要保持数据的顺序或处理重复项时,就应该使用列表。

    其次,new.iloc[new_keys[mask].index]实际上是指向原始的、未经验证的new数据框,而不是new_std数据框。

    为什么呢?因为auto_convert_columns()函数的运行目的仅仅是获取用于比较的一致类型而已。调用者仍然需要保留那些在去重处理后依然存在的原始数据及其结构。

    千万不要让某个辅助计算过程无意中成为你处理数据的最终依据。始终只修改数据的副本,然后再决定哪些内容应该保留、哪些应该删除。一旦做出了这样的决策,就要将其应用到原始数据上,这样才能确保后续步骤能够使用到真实的原始数据。

    简而言之:

    1. 最初的修改只是为了进行比较而已。

    2. 利用这种比较结果来过滤原始数据。

    3. 保留原始数据行,以便后续步骤能够继续处理它们。

    最佳实践:确保你的去重键简单且稳定(不可更改),并且在执行任何复杂的操作之前,总是先使用if existing.empty or new.empty: return new.copy()这种简洁的判断语句来避免不必要的处理。当没有数据需要处理时,这个简单的判断语句可以省去耗时的去重逻辑。

    postprocess()函数中整合所有步骤

    上述每一个步骤都是独立且易于测试的。在数据转换流程的最后阶段,这些步骤会被按顺序组合起来,形成一个完整的处理流程:

    def postprocess(df: pd.DataFrame) -> pd.DataFrame:
        """执行所有后续处理操作。」
        if df.empty:
            return df
    
        df = df.copy()
    
        print("  1. 转换数据类型...")
        df = auto_convert_columns(df)
    
        print("  2. 更改列名(从法语转换为英语)...")
        df = rename_to_english(df)
    
        print("  3> 映射分类变量值...")
        df = apply_categoricalMappings(df)
    
        print("  4> 添加衍生列..."
        df = add_derived_columns(df)
    
        print("  5> 重新排序列...)
        df = order_columns(df)
    
        print("  6> 排序数据并重置索引...")
        df = df.sort_values(["record_date", "station_code")).reset_index(drop=True)
    
        return df
    

    注意这个函数的结构:它实际上是一个包含六个编号步骤的线性处理流程,每个步骤都调用了之前定义好的纯函数。这里并没有新的逻辑代码,只有执行顺序的安排,而这种设计是刻意为之的。

    最佳实践:在每个处理阶段都要清晰地记录进度。如果一个定时任务在凌晨3点失败了,一份详细的步骤日志就能帮助你迅速诊断问题,而不会浪费一小时的时间去猜测原因。使用日志记录工具是一个不错的选择。在这个教程中,我们仍然采用了传统的控制台输出方式(print(f" 1. 转换数据类型..."))。

    第6部分:数据加载步骤

    设计逻辑

    <“加载”是整个流程的最后阶段:将处理完毕的数据传输到目标位置。常见的目标包括数据仓库、数据湖和数据库。

    在发布数据之前,数据处理流程会准备两样东西:输出文件夹本身,以及用于描述数据集的元数据文件。

    代码级说明

    create_output_dir()函数会确保输出文件夹的存在;如果该文件夹不存在,就会创建它:

    • 如果useMock的值为True,那么它会使用KAGGLE_CONFIG.mock_output_dir这个路径作为输出文件夹;否则,会使用KAGGLE_CONFIG.output_dir中指定的真实输出目录。

    • 无论哪种情况,该函数都会确保输出文件夹被创建出来。

    def create_output_dir/useMock: bool = False) -> None:
        """如果输出文件夹还不存在,就创建它。"""
        output_dir = KAGGLE_CONFIG.mock_output_dir if usemock else KAGGLE_CONFIG.output_dir
        output_dir.mkdir(parents=True, exist_ok=True)
        # exist_ok=True:这个参数使得该函数具有“幂等性”,即多次调用也会得到相同的结果。

    exist_ok=True是一个虽小但非常重要的细节。如果没有设置这个参数,Path.mkdir()会在目录已经存在的情况下抛出FileExistsError异常——而在第一次运行之后,每次执行这个函数时目录都很可能已经存在了。

    设置exist_ok=True后,创建目录的操作就会具备“幂等性”:无论调用多少次,结果都是一样的。在文件系统中,这种“可以多次安全执行”的特性正是第5部分中提到的去重逻辑所依赖的基础。

    Kaggle API遵循数据包规范。它要求在CSV文件旁边提供一个描述性强的dataset-metadata.json文件。这对于巴黎洪水数据集的搜索和被发现率来说非常重要:

    def create_metadata(df: pd.DataFrame, config: KaggleConfig) -> Dict:
        """根据DataFrame和配置信息生成Kaggle数据集的元数据。"""
        if df.empty or "record_date" not in df.columns:
            first_date = "unknown"
            last_date = "unknown"
        else:
            first_date = df["record_date"].min().strftime("%Y-%m-%d")
            last_date = df["record_date"].max().strftime("%Y-%m-%d")
    
        return {
            "title": config.title,
            "id": config.dataset_slug,
            "licenses": [{"name": config.license_name}],
            "keywords": config_keywords,
            "temporalCoverage": {"startDate": first_date, "endDate": last_date},
            "geospatialCoverage": config.geospatial_coverage,
            "updateFrequency": config.update_frequency,
        }
    

    元数据来源于KaggleConfig数据类。需要注意的是,temporalCoverage这个数值是根据数据本身计算得出的(通过df["record_date"].min()/.max()得出),而不是被硬编码设定的。每次管道运行时,元数据所对应的日期范围都会自动反映当前的实际情况,完全不需要进行任何人工设置。

    最后,publish_to_kaggle()函数会将更新后的数据集上传到Kaggle平台。它在执行过程中会调用Kaggle CLI这个外部命令工具。

    它的具体操作步骤如下:

    1. 获取当前时间,并将其转换成文本格式,例如每周更新:2026-07-13 14:30:00

    2. 构建一条Kaggle命令,该命令的内容包括:

      • 发布一个新的数据集版本

      • 使用位于KAGGLE_CONFIG.output_dir目录中的文件

      • 附加一段说明信息

      • 将目录内的所有文件压缩成ZIP文件

    3. 打印出这条命令,以便大家能清楚地看到它具体会执行哪些操作

    4. 使用subprocess.run(...)来运行这条命令

    如果上传成功,系统会输出已成功上传到Kaggle平台。

    如果上传失败,系统会打印出错误信息,并再次抛出异常,以确保错误不会被忽略。

    def publish_to_kaggle() -> None:
        """使用Kaggle CLI将更新后的数据集上传到Kaggle平台。"""
        timestamp = pdTimestamp.now().strftime("%Y-%m-%d %H:%M:%S")
        message = f"每周更新:{timestamp}"
    
        cmd = [
            "kaggle", "datasets", "version",
            "-p", str(KAGGLE_CONFIG.output_dir),
            "-m", message,
            "--dir-mode", "zip",
        ]
    
        print("正在将数据集上传到Kaggle平台...")
        print("命令内容:", " ".join(cmd))
    
        try:
            subprocess.run(cmd, check=True)
            print("已成功上传到Kaggle平台.")
        except subprocess.CalledProcessError as e:
            print(f"在上传过程中遇到错误:{e}")
            raise
    

    publish_to_kaggle函数是通过Python的subprocess模块来调用Kaggle CLI的。这是从Python管道中运行外部命令行工具的标准方法。

    有两条做法值得大家作为常规习惯来采用,而不仅仅是在使用Kaggle时:

    • cmd变量是以列表形式创建的,而不是被连接成字符串。将参数列表传递给subprocess.run可以完全避免调用shell,从而有效防止shell注入攻击的风险,同时也能正确处理那些包含空格或特殊字符的参数。有关更多详细信息,请参阅Python官方文档中关于subprocess的安全注意事项

    • 设置check=True意味着,如果Kaggle CLI执行失败并返回了非零状态码,subprocess.run会抛出CalledProcessError异常,而不会默默地退出程序。如果不设置check=True,那么上传失败的结果与成功的结果对于其他代码来说是完全相同的。

    第7部分:构建完整的处理流程

    全局测试(模拟模式)

    整个处理流程都是单独构建并经过测试的。现在,我们需要将这些各个环节串联起来,形成一个连续的功能链。首先,使用模拟数据来运行整个流程。这样并不会有任何数据被实际发布到外部。出于演示目的,一些包含人工生成数据的CSV文件会被保存到本地的`mock_output`文件夹中。

    def run_etl_pipeline/useMock: bool = True) -> pd.DataFrame:
        """运行完整的ETL处理流程:数据提取 → 数据转换 → 导出为CSV文件。"""
    
        # 第1步:加载现有数据
        if useMock:
            loaded_df = pd.DataFrame(mock_api_response)
            loaded_df['date_obs_elab'] = pd.to_datetime(loaded_df['date_obs_elab'])
            loaded_df = rename_to_english(loaded_df)
        else:
            loaded_df = load_csv(KAGGLE_CONFIG.input_csv)
    
        # 第2步:确定需要更新的数据范围
        should_update, start_date = determine_update_range(loaded_df)
        if not should_update:
            return loaded_df   # 如果数据已经是最新的,就无需进行任何操作
    
        # 第3步:提取数据
        fetched_data = fetch_all_data(start_date, useMock=useMock)
    
        # 第4步:去除重复数据
        deduped_fetched_data = remove_duplicates(loaded_df, fetched_data)
    
        # 第5步:合并数据
        new_df_english = rename_to_english(deduped_fetched_data)
        merged_historical_and_new = pd.concat([loaded_df, new_df_english], ignore_index=True)
    
        # 第6步:进行后处理
        processed_records = postprocess(merged_historical_and_new)
    
        # 第7步:导出结果文件
        create_output_dir/useMock=useMock)
        output_path = KAGGLE_CONFIG.mock_output_filename if useMock else KAGGLE_CONFIG.output_csv_path
        processed_records.to_csv(output_path, index=False, sep=",")
    
        return processed_records
    

    上述七个步骤中的每一个,都对应着你在本教程前面已经构建并测试过的函数。将这些函数串联起来其实是一个相当机械性的过程。而这正是将功能分解为职责明确的小模块所带来的好处:

    loaded_df                 = load_csv(...)                          # 1. 加载现有数据
    should_update, start_date = determine_update_range(...)            # 2. 判断是否需要更新数据
    fetched_data               = fetch_all_data(start_date)             # 3. 获取新数据
    deduped_fetched_data       = remove_duplicates(existing, new_raw)   # 4. 去除重复数据
    merged_historical_and_new  = pd.concat([existing, new_clean])       # 5. 合并数据
    processed_records          = postprocess(merged_historical_and_new) # 6. 进行后处理
    processed_records.to_csv(...)                                       # 7. 保存结果文件
    write_metadata() + publish_to_kaggle()                               # 8. 发布到Kaggle平台

    每一行的功能含义都是显而易见的,只需从左到右阅读就能理解。这正是良好代码分解结构的意义所在。

    main(): 管道编排流程

    main()充当着“指挥者”的角色。它按照正确的顺序将之前构建的所有组件连接起来,其工作原理与postprocess()类似,只不过postprocess()是在较低层级执行相同操作的。

    def main() -> None:
        """执行完整的巴黎洪水监测ETL管道流程。
    
        **数据提取阶段**
        1. 从Kaggle输入路径加载现有数据集
        2> 判断是否需要更新数据(以及需要更新的日期)
        3> 通过Hub'Eau API获取新数据
    
        **数据转换阶段**
        4> 与现有数据集进行去重处理
        5> 合并数据并进行后续处理(如类型解析、数据转换、生成新的列等)
    
        **数据写入阶段**
        6> 将更新后的CSV文件及元数据保存下来
        7> 将新版本的数据集发布到Kaggle上
    
        *如果数据集已经是最新的,程序会提前结束执行。*"
        final_dataset = run_etl_pipeline/useMock=False)
        final_dataset.to_csv(KAGGLE_CONFIG.output_csv_path, index=False)
        write_metadata(final_dataset)
    
        print("\n[最终步骤] 正在将数据发布到Kaggle...")
        publish_to_kaggle()
    
        print("正在执行运行后的验证操作")
        validate_and_analyze(final_dataset)
    
    
    if __name__ == "__main__":
        main()
    

    if __name__ == "__main__":保护机制

    if __name__ == "__main__":
        main()
    

    这是Python中最常见的编程惯用法之一,因此弄清楚它的具体作用是非常重要的。根据Python官方文档中对__main__的解释:

    • 当直接运行该文件时(例如python script.py),

      • 特殊变量__name__会被设置为"__main__"

      • 因此条件if __name__ == "__main__":会成立,

      • 从而会执行main()函数。

    • 但如果将该文件作为模块在其他地方导入(例如import script),

      • 变量__name__会变为该模块的名称,

      • 因此条件if __name__ == "__main__":不成立,

      • 所以main()函数不会自动被执行。

      最佳实践:无论何时都应该使用这种机制来保护你的程序入口点。这样,在测试单个功能或在其他脚本中重用该模块时,就可以确保不会因为简单导入该模块就触发整个管道流程(包括将数据发布到Kaggle上!)。

      第8部分:上线并切换到真实API

      以上所有步骤都是在使用模拟数据的情况下安全执行的。如果想要让管道流程直接访问真实的Hub'Eau API,需要按照以下步骤操作:

      1. 注意上面的设置useMock=False。这个设置会影响到APIConfig.usemock的值,进而影响run_etl_pipeline/useMock=False)main()函数的执行结果。

      2. 在使用数据集之前,请确保KAGGLE_CONFIG.dataset_sluginput_csv指向的是你自己的Kaggle数据集副本。只有你自己拥有的数据集,才能发布新版本。你需要同时复制这个笔记本文件这个数据集,然后更新KaggleConfig配置文件中对应的路径信息。

      3. 如果你计划在Google Colab或Kaggle笔记本之外调用publish_to_kaggle()函数,那么需要先安装并配置Kaggle CLI工具

      其他所有内容都无需进行任何更改

      运行后的验证

      validate_and_analyze()这一功能完成了整个流程的闭环控制,它不仅仅是一个可有可无的附加选项,而是一项非常重要的检查措施——该操作会在数据处理流程完成后立即执行。其输出结果是一份便于人类阅读的报告,其中包含了数据的格式、数据类型、空值数量、water_level_mm相关的统计信息、出现多少条洪水警报记录,以及各监测站点的记录数量等信息。

      这个功能并不会修改任何数据。它的存在意义在于:无论哪位人员或哪个监控系统阅读处理流程的日志输出,都能立刻判断出本周的数据是否正常。因此完全没有必要手动分析CSV文件,赶紧试试看吧!

      总结

      你刚刚构建了一个完备的ETL数据处理流程,可以反复运行而不会导致任何问题。这种数据处理框架适用于几乎所有的定期数据处理任务——只需更换API接口、字段映射方式以及数据存储目标即可。

      如果你想了解更多类似的教程,可以访问我的GitHub账号Kaggle个人主页

      感谢你的阅读!

      参考资料

相关文章

技术实践

Flutter前端系统设计:在人工智能时代,如何像资深工程师一样思考

系统设计长期以来一直被视为后端领域的问题。 如果你问一群Flutter工程师“系统设计到底意味着什么”,他们中的大多数人会提到服务器架构:负载均衡器、数据库以及微服务。 但如果你让他们设计一个分布式缓存系统或画出一个消息队列的示意图,他们会犹豫不决。而当你要求他们为社交Feed应用开发Flutter客户端时,他们就会立刻打开新文件开始编写组件代码。 这种差距确实存在,不过正在迅速缩小。 随着Flutter应用程序变得越来越复杂——它们具备了实时功能、离线支持、多平台兼容性,同时还包含需要维护的人工智能生成代码——在编写任何一个组件之前所做出的架构决策,其重要性已经与后端架构相当了。 在那些以产

阅读全文
技术实践

Flutter中的低功耗蓝牙技术:开发者手册

大多数Flutter教程都只涉及到网络调用和REST API。但一旦你需要与物理设备进行交互——比如心率监测器、智能灯泡、健身追踪器、工业传感器,或者你自己定制的硬件设备——你就不得不离开HTTP这个“舒适的环境”,转而使用蓝牙低功耗技术。 本指南会教你如何在Flutter中正确且全面地实现这些功能。 移动设备上的蓝牙功能其实相当复杂。Android和iOS之间的权限设置有所不同,即使是同一款Android系统的不同版本,权限要求也会存在差异。蓝牙连接的生命周期包含许多状态,服务与特征的数据模型也会让新手感到困惑,而字节级的数据编码方式几乎会让每个人在初次尝试时遇到麻烦。 flutter_bl

阅读全文
技术实践

如何在现代API中实施“以隐私保护为导向的设计理念”——开发人员的实用指南

作为软件开发人员,我们通常被教导要优先考虑速度、性能和正常运行时间等因素。在构建API时,我们的核心目标就是确保数据能够顺利地从A点传输到B点。 然而,全球范围内的数据隐私法规正在日益严格,用户也越来越关注自己的数字足迹。将隐私问题视为“事后才需要处理的法律事项”,或者认为可以在生产环境中再解决这些问题,已经不再是一种可持续的做法。 这时,《设计即隐私》这一理念就派上了用场。 “设计即隐私”是一个框架,旨在将隐私保护措施主动融入到工程开发的整个生命周期中。这意味着,你的系统架构应该默认就能保护用户数据。 在这份全面的指南中,我们将通过现代的工程模式、代码设计理念以及有针对性的数据库结构,来讲解

阅读全文
技术实践

Firestore数据建模指南:内嵌文档与引用方式的选择(结合一个博客案例进行分析)

当开发人员从关系型数据库世界(如MySQL、PostgreSQL)转向Firebase这种NoSQL文档数据库时,他们往往会沿用原有的习惯,试图在新的系统中复制表结构、外键以及关联操作。 其结果是什么呢?复杂的查询语句、飞涨的读取成本,以及那种在添加少量功能后就变得难以维护的数据库结构。 要理解Firestore的工作原理,我们首先需要对比它所基于的关系型模型。只有弄清楚SQL是如何处理数据的,才能清楚地看到Firestore与它的区别所在,从而了解如何正确地构建NoSQL数据结构。 在本指南中,我们将介绍NoSQL的设计原则、数据嵌入与引用机制,以及各种类型的关系建模方法(1-1、1-N、N

阅读全文