Towards Data Science

I Built My Second ETL Pipeline. This Time, I Started Thinking Like a Data Engineer

6.9内容质量

TL;DR · AI 摘要

I Built My Second ETL Pipeline. This Time, I Started Thinking Like a Data Engineer Towards Data Science Data Engineering...

核心要点

  • 主题聚焦:I Built My Second ETL Pipeline. This Time, I Sta
  • 来源:Towards Data Science,建议结合原文判断细节。
  • AI 分析暂不可用,本条为保底评分与摘要。
#AI#编程#后端#云计算#安全
打开原文

我构建了第二个ETL管道。这一次,我开始像数据工程师一样思考 | Towards Data Science

数据工程

我构建了第二个ETL管道。这一次,我开始像数据工程师一样思考

使用Python、Docker、PostgreSQL和Kestra构建生产级RSS管道

Ibrahim Salami

2026年7月10日

18分钟阅读

分享

AI生成的架构图

几个月前,我决定想要从数据分析师转型为数据工程师。

和许多刚入门的人一样,我被需要学习的事项数量所压倒。数据仓库、编排工具、分布式处理、流处理系统、云平台、基础设施。这个列表似乎永无止境。

我没有尝试一次性学习所有内容,而是采取了不同的方法。

我制定了一个为期12个月的自学路线图,围绕一个简单的理念构建:通过构建来学习。

而不是在教程之间跳跃,我会构建一系列小型项目。每个项目将引入几个新概念,同时强化之前学到的内容。目标不是构建最复杂的应用程序,而是通过逐步解决问题,培养像工程师一样思考的习惯。

这段旅程的第一个项目是一个GitHub ETL管道。它最初只是一个简单的Python脚本,用于获取仓库数据并导出到CSV文件。随着学习的深入,我逐步改进了它。我用SQLite替换了CSV文件,使管道具备幂等性以防止重复数据,并最终通过GitHub Actions实现了自动化。

当我完成这个项目时,我意识到了一些在开始时并不明显的事情。

构建ETL逻辑实际上只是容易的部分。

更困难的问题出现在我停止思考一次性运行的脚本,开始思考需要反复运行且无需人工监控的系统时。

  • 应该如何安排其执行时间?
  • 如果它在执行中途失败了怎么办?
  • 重试逻辑应该放在哪里?
  • 如何打包应用程序以确保它在任何地方都能以相同方式运行?

这些问题让我对数据工程的理解远超过解析JSON或编写SQL。

第一个项目给我带来了另一个挑战。

GitHub Actions对于自动化小型ETL管道效果不错,但我想了解当使用专门为数据工程构建的工作流编排器时会发生什么变化。我想学习工程师如何将编排与执行分离,容器化工作负载如何融入其中,以及一个面向生产的管道在小规模下实际是什么样子。

因此,我的第二个项目决定构建一个自动化的RSS数据摄入管道。

表面上看,这是一个相当简单的应用。它从RSS源获取文章,将其解析为结构化对象,并存储到PostgreSQL中。

但目标从来不是构建一个RSS阅读器。

目标是探索将Python脚本转化为可靠数据管道所需的工程决策。

在本文中,我将逐步介绍这些决策,我在过程中犯下的错误,以及使用Kestra构建第一个管道时学到的经验教训。

为什么要构建另一个ETL管道?

在完成第一个ETL项目后,我考虑转向完全不同的方向。

也许是数据仓库项目。也许是Apache Spark。也许是带有更复杂转换层的API。

相反,我构建了...另一个ETL管道。

起初,这听起来可能像是倒退了一步。

毕竟,我之前已经构建了一个提取管道,使其具备幂等性,并通过 GitHub Actions 进行了调度。为什么要重复同样的工作?

因为我并非在学习一个新的数据集,而是在学习一种新的思维方式。

我第一个项目中的一个教训一直铭记于心。

编写 ETL 逻辑本身并不是困难的部分。困难的部分在于围绕它的所有其他问题。

  • 管道应该如何执行?
  • 它应该如何从失败中恢复?
  • 如何打包它,以确保在任何机器上都能一致运行?
  • 调度应该属于哪里?

也许最重要的是:

应用程序的责任应该在哪里终止,编排层的责任又应该从哪里开始?

这些问题与 RSS 订阅源或 GitHub 仓库关系不大。它们是工程问题,我意识到几乎可以用任何数据源来探索这些问题。

这就是我选择 RSS 订阅源的原因。

不是因为 RSS 特别令人兴奋,而是因为它刻意保持简单。

提取逻辑只需要几行 Python 代码。这意味着我可以花更少时间担心业务逻辑,而将更多时间用于思考架构。

出于同样的原因,我决定在这个项目中不再使用 GitHub Actions。

GitHub Actions 是一个很好的调度入门工具。它向我展示了如何自动化工作流,并让我第一次体验了无需人工干预即可运行 ETL 管道。

但 GitHub Actions 并不是专门为数据工作流编排而设计的。

我想了解使用专为数据管道构建的工具时会发生哪些变化。

这正是我转向 Kestra 的原因。

与其问“如何每小时运行这个 Python 脚本?”,我发现自己在思考另一组问题。

  • 重试应该如何配置?
  • 工作流执行是如何被跟踪的?
  • 环境变量应该如何传递到容器中?
  • 失败的执行会呈现什么样子?
  • 如何将应用程序与运行它的基础设施分离?

这些问题正是我想探索的。

通过选择一个简单的 ETL 管道,我可以专注于工程决策,而不会被复杂的业务逻辑分散注意力。

回顾过去,我认为这是正确的决定。

这个项目之所以有趣,并不是因为它处理了 RSS 订阅源。

它之所以有趣,是因为构建它迫使我去思考可靠性、可重复性和编排,而这些是我第一个项目从未涉及过的。

第一个架构决策:在 Kestra 之前使用 Docker

在项目定义之后,我的第一反应是直接跳入 Kestra。

毕竟,编排是我选择这个项目的主要原因之一。为什么不从这里开始?

相反,我做了一件后来证明能节省我很多麻烦的事。

我完全忽略了 Kestra。

这听起来可能有些违反直觉,但我想在确认核心应用真正运行之前,避免引入多个变动因素。

因此,我以分层的方式构建了这个项目。

首先,我编写了 Python ETL 代码。

它的职责被刻意设计得非常简单。获取 RSS 订阅源,将每条条目解析为 Article 对象,并将结果保存到 PostgreSQL。仅此而已。

python
feed = feedparser.parse(RSS_URL)

articles = parse_feed(feed)

save_articles(articles)

ETL 的工作内容基本上就是这些。我刻意将应用程序保持得较为简单,因为本项目的核心并不在于转换逻辑,而是在它周围发生的一切。

当这个流程能够稳定运行后,我将注意力转向了数据库。

为了确保重复执行的安全性,我使用 PostgreSQL 的 ON CONFLICT DO NOTHING 使插入操作具备幂等性。这样即使多次运行管道,也不会产生重复的行。

code
INSERT INTO articles (...)
VALUES (...)
ON CONFLICT (id) DO NOTHING;

这一行代码让管道可以安全地重复执行。无论 Kestra 是运行一次还是运行一百次工作流,PostgreSQL 都会负责防止重复记录的产生。

只有在 ETL 和数据库协同工作后,我才引入了 Docker。

这个决定彻底改变了我对项目的思考方式。

最初,Docker 看起来只是另一个需要学习的工具。

但项目结束时,我意识到它已经变得比预期更加重要。

它成为了执行的基本单元。

code
FROM python:3.13-slim

WORKDIR /app

COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

COPY . .

CMD ["python", "fetch_rss.py"]

以这种方式打包 ETL 使应用程序变得自包含。我不再需要让 Kestra 理解我的 Python 项目,而是可以简单地要求它运行一个已经知道如何执行管道的容器。

我的思维方式也发生了转变:不再思考“Kestra 会运行我的 Python 脚本”,而是开始思考“Kestra 会运行我的 Docker 镜像”。

这个区别看似细微,却彻底改变了你的应用程序与编排层之间的关系。

当 ETL 被打包进容器后,无论它是在我的笔记本电脑上运行、在 Kestra 内部运行,还是在其他机器上运行,都无关紧要。

运行时环境始终是一致的。

这种一致性给了我之前没有的东西:信心。

在引入 Kestra 之前,我可以手动运行容器并验证它是否能获取 RSS 源、连接到 PostgreSQL 并持久化预期的记录。

如果出现故障,我知道问题不会隐藏在另一层编排之下。

这个验证步骤后来证明极其有价值。

在开发过程中,我遇到了网络、容器配置和工作流执行的问题。由于 Docker 镜像已经过独立验证,我可以立即排除 ETL 本身的问题,专注于编排层。

这使调试变得显著更容易。

回顾整个项目,我认为这是学到的最重要的教训之一。

人们很容易倾向于尽快将所有组件连接起来并希望一切都能正常工作。

更好的方法是在引入下一层之前先验证每一层。

在这个项目中,验证顺序如下:

  • 验证 Python ETL
  • 验证 PostgreSQL 持久化
  • 验证 Docker 镜像
  • 最后,让 Kestra 编排一个我已经信任的容器

每一层都建立在前一层的基础上。

当 Kestra 进入项目时,我并不是同时要解决 Python、PostgreSQL、Docker 和编排的问题。

我只专注于解决一个问题。

这种渐进式方法让整个项目变得更容易管理,而且这种工作流程我可能会在未来的数据工程项目中继续使用。

最终被证明错误的假设

有了可用的 Docker 镜像后,我确信最难的部分已经过去了。

  • 我有一个能正常运行的 Python ETL 流程。
  • 我在 Docker 中运行了 PostgreSQL。
  • 我已经验证过容器可以抓取 RSS 文章并保存到数据库中。

现在 Kestra 只需要运行它即可。

或至少我这么认为。

我最初的假设很简单。

Kestra 会指向我的 Python 文件,执行脚本,一切都会和从命令行运行时完全一致。

事实并非如此。

我很快发现了一个我之前没有充分理解的重要区别。

Kestra 是一个编排器。

它不负责构建 Python 环境或管理应用依赖。它的职责是决定任务何时以及如何运行。

这个认知改变了整个项目的方向。

我不再把 Kestra 看作另一个执行 Python 代码的地方,而是开始把它视为负责编排已有工作负载的层级。

那个工作负载就是我的 Docker 镜像。

一旦我完成了这个认知转变,架构变得清晰得多。

  • ETL 变成了一个自包含的应用程序。
  • Docker 成为了部署产物。
  • Kestra 成为了编排器。

每一层都有明确的职责,它们不需要了解其他层的内部实现方式。

有趣的是,达到这个阶段的过程并不完全直接。

我最初的直觉是想办法让 Kestra 直接执行 Python 项目。这种方法听起来更简单,但随着我不断尝试,我逐渐意识到自己在要求编排器承担本应由应用程序本身提供的责任。

一旦我接受 Docker 作为执行产物,整个流程变得简单得多。

一些方法起初看起来很有希望,但引入了不必要的复杂性。其他方法虽然可行,却与 Kestra 推荐的容器化工作负载执行方式不一致。

最终,我找到了一个令人惊讶地简单的流程。

我不再试图教 Kestra 如何运行 Python,而是让 Kestra 做它最擅长的事情。

它启动一个容器。

code
tasks:
  - id: run_etl
    type: io.kestra.plugin.scripts.shell.Commands

    containerImage: rss-pipeline-etl:latest

    taskRunner:
      type: io.kestra.plugin.scripts.runner.docker.Docker

    commands:
      - python /app/fetch_rss.py

现在再看这个流程,最引人注目的不是 YAML 的复杂程度。而是 Kestra 实际上需要了解我的应用内容非常有限。它唯一的职责就是启动一个已经知道如何执行 ETL 的容器。

在这个容器内部,我的应用已经完全清楚该做什么。

这个小小的架构调整不仅解决了即时的执行问题。

它还强化了一个在我学习过程中反复出现的理念。

优秀的工程往往不是关于添加更多层级。

而是让每一层只承担单一职责,并让它出色地完成这个任务。

  • Python 不应该关心编排。
  • Kestra 不应该关心 Python 依赖。
  • Docker 不需要了解任何关于 RSS 消息的内容。

每个组件解决不同的问题。

一旦我停止要求一个工具解决所有问题,整个系统变得更容易理解。

回顾过去,这可能是整个项目中最大的思维转变。

我不仅学会了如何使用 Kestra。

我还理解了编排真正的含义。

一旦流水线能够自动运行,一切都会改变

到目前为止,我一直是手动运行这个流水线。

如果某个环节失败了,我需要坐在电脑前。我可以阅读错误信息,进行修改,然后再次尝试。

一旦流水线开始自动运行,这种安全网就消失了。

我在 Kestra 中配置的第一个功能就是一个简单的每小时定时任务。

code
triggers:
  - id: hourly_schedule
    type: io.kestra.plugin.core.trigger.Schedule
    cron: "0 * * * *"

从纸面来看,这只是一个 cron 表达式。

但在实际应用中,它代表了更深远的转变。

流水线不再依赖我记住去手动运行它。

每个小时,Kestra 都会启动一个新的执行,启动 ETL 容器,并处理 RSS 源中的最新文章。

这立即引发了一个新问题。

如果这些执行中的某一次失败了怎么办?

在开发过程中,我故意引入了故障来验证这个问题。我将流水线指向了一个无效的数据库主机,观察发生了什么。

第一次执行确实如预期般失败了。

更重要的是,它并没有就此停止。

由于工作流配置了重试机制,Kestra 在短暂延迟后自动重新尝试了执行。

code
retry:
  type: constant
  maxAttempts: 3
  interval: PT30S

这成为我项目中最喜欢的时刻之一。

当我修正了配置后,工作流在不修改应用程序本身的前提下成功完成。

这并不是因为重试机制特别复杂,而是因为它突显了责任划分的另一个层面。

ETL 不应该决定自己是否值得再次尝试。

这是编排层需要考虑的问题。

通过将重试逻辑移至 Kestra,Python 应用程序可以专注于单一职责:处理源数据并持久化结果。

编排层负责处理系统的韧性。

定时任务引入了另一个挑战,这在第一次 ETL 项目中我就已经遇到过。

重复执行意味着重复尝试写入数据。

如果同一篇文章在多个小时运行中出现,流水线不应该重复插入两次。

幸运的是,我之前已经解决过类似的问题。

数据库层通过使用 PostgreSQL 的 ON CONFLICT DO NOTHING 设计为幂等性。

这意味着每次执行都可以安全地尝试插入相同记录,而不会产生重复数据。

定时任务、重试机制和幂等性写入的结合,使流水线变得更加宽容。

如果某次运行失败,Kestra 可以重试。

如果成功的重试遇到已经写入的数据,PostgreSQL 会直接忽略重复项。

每一层都不需要知道另一层在做什么。

它们各自处理自己的职责。

最后的改进是可见性。

在开发初期,我的日志看起来和你预期的调试阶段项目完全一样。

我直接将整个 RSS 对象打印到控制台,只是为了确认解析器是否正常工作。

这并不美观,但确实完成了它的目的。

随着流水线变得越来越稳定,这些调试语句的实用性逐渐降低。

我将它们替换为描述执行过程的日志,而不是直接输出原始数据。

之前

code
print(feed.entries[0])

之后

code
print("=== RSS PIPELINE START ===")

first = feed.entries[0]

print("First entry preview:")
print(f"Title: {first.get('title')}")
print(f"Link: {first.get('link')}")

print(f"Feed title: {feed.feed.title}") print(f"Fetched articles: {len(articles)}") print(f"Saved {len(articles)} articles to the database.")

print("=== RSS PIPELINE END ===")

code

现在每次运行都会讲述一个简单的故事。

=== RSS PIPELINE START ===

第一条条目预览: 标题:Christian Ledermann: 从 mypy 迁移到 ty 和 pyrefly 链接:https://dev.to/...

Feed 标题:Planet Python 获取文章数:25 已将 25 篇文章保存到数据库。

=== RSS PIPELINE END ===

code

这些改动并没有让应用变得更聪明,而是让其更容易理解。这是一个重要的区别。

良好的可观测性不在于产生更多的日志,而在于产生正确的日志。

项目结束时,我意识到一个有趣的事情:Python 代码并没有显著增长,大部分工作都发生在代码之外。

回顾来看,大部分工程努力并不是用于让 ETL 更聪明,而是用于让其更可靠。

- 定时任务确保了它可以在没有人工干预的情况下运行。
- 重试机制帮助其从临时故障中恢复。
- 幂等性保护数据库免受重复写入的影响。
- 日志让每次执行都更容易理解。

单独来看,这些改动都不算特别复杂。但结合起来,它们将一个简单的 Python 脚本转变成了一个行为更像生产系统的应用。

## 最终架构

项目结束时,架构最终演进成一个令人惊讶地简单的形态。

看着最终的架构,人们很容易认为项目一开始就朝着这个方向发展。

事实并非如此。

每一层只有在前一层得到验证后才会被添加。

- 首先是 ETL。
- 然后是 PostgreSQL。
- 接着是 Docker。
- 最后是 Kestra。

这个顺序非常重要。

因为每个组件都已经被独立测试过,我从未发现自己需要同时调试 Python、Docker、PostgreSQL 和 Kestra。每个决策都减少了未知因素,而不是增加它们。

更重要的是,每个组件最终都明确了自身的职责。

- Python 负责获取、解析和持久化 RSS 文章。
- PostgreSQL 负责安全地存储数据并防止重复。
- Docker 提供一致的执行环境。
- Kestra 决定工作负载何时运行,以及在失败时应该发生什么。

这些组件都没有试图去做彼此的工作。讽刺的是,这正是让最终系统比我预期的更简单的原因。

## 这个项目如何改变了我的思维方式

当我刚开始学习数据工程时,我假设困难的部分会是编写 ETL 代码。

这正是大多数入门教程关注的重点。

- 你学习如何调用 API。
- 你转换数据。
- 你将其保存到某个地方。
重复这个过程。

这些是宝贵的技能,但它们只是整个图景的一部分。这个项目教会我,工程工作真正开始于脚本能够运行之后。

一旦一个管道被期望每小时运行一次、能够从临时故障中恢复、避免重复数据并产生解释发生了什么的日志,问题就会变得有趣得多。

你不再思考单个函数,而是思考系统。

对我而言,最大的思维转变之一是理解执行与编排之间的区别。

起初,这些概念似乎几乎可以互换。现在它们却感觉完全分开。

Python 应用应专注于业务逻辑,编排器应专注于应用何时、何地以及如何运行。

将这些职责分开使整个项目更容易理解。

这也改变了我对 Docker 的看法。

在该项目之前,我主要将 Docker 看作一种打包应用的方式。

现在我将其视为一种部署产物。

一旦 ETL 被打包到容器中并独立验证,我就无需再担心它在 Kestra 中是否会表现不同。

这种信心最终证明是容器化应用最大的优势之一。

不过,最重要的教训与 Kestra 或 Docker 无关。这是关于渐进构建的价值。

每个重大决策都遵循相同的模式:

- 构建最小可行的解决方案
- 验证其有效性
- 之后才引入下一层

这种方法使项目感觉没有一开始就试图连接所有组件那样令人望而生畏。回顾过去,我认为这将是我未来每个项目都会坚持的教训,无论使用何种技术。

## 展望未来

这个 RSS 管道只是我数据工程学习旅程中的第二个项目。与我的第一个 ETL 管道相比,代码本身并没有显著更复杂。改变的是我解决问题的方式。

我不再问"如何编写这个脚本",而是开始思考:

- 这个责任应该放在哪里?
- 如果它失败了会发生什么?
- 能否在不担心重复数据的情况下反复运行?
- 能否在我不监控时可靠运行?

这些问题促使我从单纯的 Python 编码者思维,转向系统设计者的思维方式。我怀疑这正是构建项目的真正价值所在。

每个项目都会教会你新工具。但最好的项目会慢慢改变你的思维方式。这个项目确实如此。

我确信下一个项目将挑战完全不同的假设集。说实话,我期待发现这些新假设。

这是我在持续记录从系统分析师向数据工程师转型系列文章的一部分。如果你一直关注,非常感谢。

在 LinkedIn、YouTube 和 Twitter 上关注我。

作者:Ibrahim Salami

查看 Ibrahim Salami 的所有文章

数据工程师

,

Docker

ETL

ETL 管道

Python

分享本文

- 在 Facebook 上分享
- 在 LinkedIn 上分享
- 在 X 上分享

Towards Data Science 是一个社区出版物。提交你的见解以触达全球受众,并通过 TDS 作者支付计划获得收益。

更新为你的实际投稿链接

为 TDS 写作

✦ 结束 CTA ✦