如果你已经不再需要使用原有的定时任务调度工具了,应该该怎么办呢?
大多数开发者在开始自动化工作的过程中都会采取类似的方法。他们会编写一些脚本来完成某些有用的任务,比如从API中获取数据、调整一批图片的大小,或者发送报告邮件,然后将这些脚本安排在每天早上自动运行。 这样做之后,他们会在自己的crontab文件中添加相应的指令,这样就能感受到自己对自动化流程的控制力了。现在,在他们睡觉的时候,电脑会自动执行这些脚本。 有一段时间里,这样的安排确实足够用了。但后来就不再适用了。 也许某个备份脚本在凌晨3点悄悄地失败了,而你直到当天晚些时候需要使用该脚本时才发现了这个问题;又或者你编写了一个在终端环境中运行得非常顺利的脚本,但却发现在用crontab安排它执行时却会
障碍一:步骤之间的依赖关系
你的晨间例行流程原本只包含一个脚本,现在却发展成了三个脚本。以一个包含三个脚本的ETL处理流程为例:
extract.py脚本从API中获取昨天的订单数据。
transform.py脚本对数据进行清洗并计算总数。
load.py脚本将处理结果写入分析数据库。
每个步骤都依赖于前一个步骤的结果。使用Cron调度工具时,常见的做法是随意设定各个脚本的执行时间:
0 2 * * * python extract.py
0 3 * * * python transform.py
0 4 * * * python load.py
你希望extract.py能在一小时内完成执行,这样transform.py才能有可处理的数据。但如果在某个特定日子里API响应速度变慢,extract.py需要花费70分钟才能完成任务,那么transform.py就会使用过时或缺失的数据进行运算,从而产生错误的结果。Cron工具并不具备“只有当A步骤成功完成后才执行B步骤”的功能,它只关注实际运行时间。
障碍二:故障处理与重试机制
网络故障、API返回503错误代码、数据库连接中断……一个可靠的自动化流程必须能够检测到这些故障并自动重试。也许应该重试三次,或者每次重试之间间隔一定的时间,这样就不会过度负担那些已经处于高负荷状态的服务。
使用Cron工具时,重试逻辑需要由开发者自行编写。这意味着你得在每个脚本中都添加相应的重试代码:使用try/except语句、设置延迟时间、使用计数器以及标志文件来判断前一个步骤是否已完成。如果同时处理多个这样的流程,那么最终就会得到一个结构复杂、容易出错且缺乏文档支持的自动化系统。
障碍三:流程的可监控性
你有没有想过:昨晚这些脚本是否真的运行了?哪些步骤成功完成了?每个步骤花费了多少时间?是extract.py步骤出现问题,还是load.py步骤出了故障?
使用Cron工具时,要得到这些问题的答案其实非常困难。首先,你得去找那些记录脚本执行情况的日志文件——而这些文件的存放位置往往并不固定。
有时候,日志文件可能保存在/var/log/syslog目录中,而有时候则位于/var/log/cron目录中;此外,只有当你明确指定了脚本的输出路径时,这些输出信息才会被保存下来。因此,你不得不仔细查阅系统日志,才能确认流程是否已经运行完毕,然后再寻找那些记录了脚本执行结果的文件。
当某个流程失败时,要找出故障原因更是困难重重。因为你必须仔细分析多次运行产生的日志文件,试图判断哪个时间戳对应着最后一次成功的执行过程,以及问题出在什么地方。通常情况下,你只能依靠非零的退出代码来获取线索。由于没有仪表盘、没有运行历史记录,也没有任何提示信息,因此发现故障往往需要花费很长时间。
默认情况下,系统失败时不会发出任何警报,而这正是后台进程最危险的特性之一。只有当下游环节发现数据更新停止了,你才会意识到你的自动化流程出现了问题——而那时可能已经有一周的时间过去了。
第四部分:数据补录与重新执行流程
当你的分析数据库已经运行了两个月之后,你发现transform.py中存在一个会导致总数计算错误的漏洞。你已经修复了这段代码,但现在需要为这两个月中的每一天重新执行相应的处理流程,以便处理那些当时未被处理的数据。
这种数据补录操作使用cron来调度时,会变得非常麻烦。因为cron只会在“当前时间”执行任务,它并没有“将任务按特定顺序执行”的功能。所以你不得不编写另一个包含日期循环的临时脚本,然后祈祷这个脚本是幂等的,并且要不断地监控它的运行情况。
请注意,在这四种情况下,你都在重复去做那些本来就可以通过某些工具很好地解决的问题。而这类工具其实就是工作流编排工具。
工作流编排能否解决这些问题?
工作流编排工具用于管理工作流——这些工作流是由具有特定关系、触发条件以及错误处理机制的任务组成的集合,同时还能提供跟踪执行结果的功能。
虽然市面上有很多类似的解决方案,比如Airflow、Dagster、Prefect和Temporal,但我们将重点介绍Kestra这款开源工具。Kestra的独特之处在于它允许开发人员通过一种与任何编程语言和基础设施兼容的声明性方式来运行、监控和管理工作流,无论这些基础设施是公共网络、私有网络,还是隔离网络。
Kestra还能帮助企业确保他们能够获得所需的数据分析结果,并满足各种监管要求。由于提供了超过1,600种连接器,你可以使用它来构建几乎任何类型的数据处理、基础设施集成或人工智能工作流。
在本教程中,我们将继续以cron调度器为例,构建一个简单的数据提取、转换和加载工作流。这个工作流会使用类似cron的调度机制,同时也会弥补cron的一些缺陷,比如执行顺序问题和错误处理机制。
Kestra入门
Kestra的出现正是为了解决编写Python代码来实现复杂工作流编排所带来的种种麻烦。使用Kestra时,你的工作流只需用YAML格式进行描述即可——只需用简洁的声明性语法说明哪些任务需要执行,然后将文件保存到Git仓库中,像普通代码一样进行部署。
你不需要编写复杂的用户界面逻辑,也不需要进行隐藏的状态管理。所有的配置信息都只是简单的文本,方便查看和对比修改前后的差异。你看到的YAML文件内容,就是最终会实际执行的工作流内容。
图1中的示例YAML文件展示了一个将NoSQL数据导入到分析数据库中的场景。
id: cassandra-to-bigquery
namespace: company.team
tasks:
- id: query_cassandra
type: io.kestra.plugin.cassandra.Query
session:
endpoints:
- hostname: localhost
port: 9042
localDatacenter: datacenter1
cql: |
SELECT salary_id, work_year, experience_level, employment_type,
job_title, salary, salary_currency, salary_in_usd, employee_residence,
remote_ratio, company_location, company_size
FROM test.salary
fetchType: STORE
- id: write_to_csv
type: io.kestra.plugin.serdes.csv.IonToCsv
from: "{{ outputs.query_cassandra.uri }}"
- id: load_bigquery
type: io.kestra.plugin.gcp.bigquery.Load
from: "{{ outputs.write_to_csv.uri }}"
destinationTable: my_project.my_dataset.my_table
serviceAccount: "{{ secret('GCP_SERVICE_ACCOUNT_JSON') }}"
projectId: my_project
format: CSV
csvOptions:
fieldDelimiter:>,
skipLeadingRows: 1
图1:Cassandra与BigQuery的示例
即使对Kestra了解不多,这个YAML文件也非常简单且容易理解。在本文的后半部分,我们将创建一个基本的ETL流程,并解释id和type等字段的含义。但首先,让我们在本地机器上运行一个Kestra实例。
如何安装Kestra
Kestra既可作为开源平台使用(采用Apache 2.0许可证),也可作为企业级解决方案购买,后者会提供更多功能和产品支持。在本教程中,我们将使用Docker在本地运行最新版本的Kestra。
要启动Kestra,请执行以下Docker命令:
docker run --pull=always --rm -it -p 8080:8080
--user=root \
--name kestra \
-v kestra_data:/app/storage \
-v kestra_db:/app/data \
-v /var/run/docker.sock:/var/run/docker.sock \
-v /tmp:/tmp \
kestra/kestra:latest server local
对于Windows和Linux等其他平台,请查看相关说明。
容器启动后,请访问http://localhost:8080进入Kestra用户界面。欢迎页面会提示您创建管理员账户,请完成相应的操作。
设置完成后,点击左侧面板中的“Flows”选项卡,然后点击页面右上角的“Create”按钮。系统会使用示例模板生成一个新的流程,如下图所示:
图2:显示新流程模板的界面
Kestra中的流程管理
在Kestra中,您可以通过“Flows”来定义工作流逻辑。这些流程可以通过UI中的YAML语法创建,也可以通过UI中的无代码编辑器生成,或者通过API进行编程式配置。在本教程中,我们将使用YAML格式来创建流程。
从上图可以看出,系统已经预先创建了一个示例流程供您开始使用。
_id、namespace和tasks是三个必备字段,它们用于在Kestra环境中识别特定的流程以及该流程需要执行的操作。每个流程都属于某个命名空间。命名空间类似于文件系统中的文件夹,用于对流程进行分类和管理。需要注意的是,一旦创建了流程,其所属的命名空间是无法更改的。
如何使用Kestra
步骤1:按计划执行简单任务
首先,让我们删除现有的示例流程,然后替换为以下内容:
id: morning_report
namespace: tutorial
tasks:
- id: say_hello
type: io.kestra.plugin.core.log.Log
message: "早上好——该流程在 {{ executionstartDate }} 运行。"
这个示例包含了所需的 id 和 namespace,同时还有一个用于记录信息的 tasks 字段。其中显示的消息使用了 Pebble 表达式(即位于 {{ }} 括号内的代码),来显示流程执行的起始日期。
Pebble 表达式可用于在流程中动态设置数值。在这个示例中,执行开始的日期会被插入到字符串中。
要运行这个流程,首先需要点击页面右上角的“保存”按钮将其保存下来。保存完成后,再点击“播放”按钮,此时会出现如图 3 所示的执行选项页面:
图 3:执行流程选项
如果这个工作流需要输入文件名或 URL 等信息,你可以在这里手动输入这些值来测试流程。该对话框还提供了 curl 命令,如果你想通过 API 而不是 UI 来运行流程,也可以使用这个命令。点击“执行”按钮即可运行流程:
图 4:流程执行日志
这个流程非常简单,它只是将一条信息记录到了日志文件中。如果流程在运行过程中出现了错误或警告,你可以在这个页面上看到详细的执行日志。
现在我们已经创建了第一个任务,接下来让我们使用 cron 表达式来安排它的执行时间。点击页面顶部的“编辑流程”按钮,将以下内容添加到流程中。
接着,在流程中添加触发器部分:
triggers:
- id: every_minute
type: io.kestra.plugin.core.trigger.Schedule
cron: "* * * * *"
运行这个流程。几分钟后,点击左侧导航栏中的“执行”选项卡,查看执行历史记录。
图 5:流程执行页面
现在,我们的流程就类似于一个单独的 cron 作业了。这个流程包含了一个触发器块,其中的 Schedule 触发器的 cron 表达式采用了你已经熟悉的那种五字段格式。这一点非常重要:你并没有抛弃自己之前学到的知识,而是将这些知识应用到了一个可以根据需求变化而进行扩展的系统当中。
请注意,即使在这个简单的 Kestra 示例中,你也获得了 cron 所没有提供的功能。每次这个流程运行时,它的执行情况都会被记录下来,包括时间戳、执行时长和状态等信息,所有这些日志都可以在用户界面中查看。这就是我们在还没有开始做任何复杂的事情之前就已经实现的“可见性”这一目标。
步骤2:实际操作与依赖关系
现在,让我们用包含提取、转换和加载三个步骤的任务来替换原来的简单任务,并让调度系统来确保这些步骤按照正确的顺序执行,而不是通过调整执行时间来实现。
id: csv_to_parquet
namespace: company.team
description: 下载订单CSV文件,使用Python脚本对其进行转换,然后将结果写入Parquet格式的文件中。
tasks:
# 将公共CSV文件下载到Kestra的内部存储系统中
- id: download_csv
type: io.kestra.plugin.core.http.Download
uri: https://huggingface.co/datasets/kestra/datasets/raw/main/csv/orders.csv
# 使用简单的Python脚本对CSV文件进行转换,并将其保存为Parquet格式
- id: transform_to_parquet
type: io.kestra.plugin.scripts.python.Script
containerImage: ghcr.io/kestra-io/pydata:latest
inputFiles:
input.csv: "{{ outputs.download_csv.uri }}"
outputFiles:
- orders.parquet
script: |
import pandas as pd
# 读取下载到的CSV文件
df = pd.read_csv("input.csv")
# --- 简单的转换操作 ---
# 确保数据类型为数值型,并添加一个计算列
df["total"] = df["quantity"] * df["price"]
# 作为示例,只保留总金额大于某个阈值的订单
df = df[df["total"] > 0]
print(f"转换后的记录数:{len(df)}")
# 将结果写入Parquet文件中
df.to_parquet("orders.parquet", index=False)
# 记录Parquet文件的生成情况
- id: log_output
type: io.kestra.plugin.core.log.Log
message: "Parquet文件已创建:{{ outputs.transform_to_parquet.outputFiles['orders.parquet'] }}"
# 将生成的Parquet文件作为可下载的输出结果提供。
# 文件类型的输出结果会显示在执行的“概览”页面中,并提供一个下载按钮。
outputs:
- id: parquet_file
type: FILE
value: "{{ outputs.transform_to_parquet.outputFiles['orders.parquet'] }}"
保存这个流程配置,然后执行它。
图6:执行结果
发生了两件事。首先,按顺序列出的任务会依次执行。只有当`download_csv`任务成功完成之后,`transform_to_parquet`任务才会开始执行;而`log_output`任务则会在`transform_to_parquet`完成后才启动。如果某个任务失败了,整个执行流程就会停止,因此`transform_to_parquet`任务也不会处理已经过时的数据。这样就解决了依赖关系带来的问题,也不需要再担心执行时间的问题了。
其次,请注意流程配置中的`{{ outputs.download_csv.uri }}`这个表达式。通过这样的表达式,任务可以将数据和元数据传递给后续的任务。这种机制使得一系列脚本能够真正构成一个连贯的处理流程。
步骤3:通过重试机制应对失败
现在,我们来考虑这样一种情况:`download_csv`任务遇到了网络问题,导致流程无法下载到最新的数据。我们可以通过在`download_csv`任务中添加重试机制,使其具备恢复能力:
- id: download_csv
type: io.kestra.plugin.core.http.Download
uri: https://huggingface.co/datasets/kestra/datasets/raw/main/csv/orders.csv
retry:
type: constant
maxAttempts: 5
interval: PT10S
这就是整个重试策略。如果下载失败,Kestra会等待并重新尝试最多5次,每次尝试之间会有10秒的延迟(PT10S表示“10秒钟”)。这个机制不需要使用任何计数器、睡眠函数或标志文件。错误处理部分仅由三行代码完成,读起来就像一句话一样简单。
要测试这一机制,只需将orders.csv文件中的“s”字符删除,然后重新运行流程,你就会看到系统会尝试重试的执行过程。
图7:执行过程中显示的重试情况
如果所有尝试都失败了,你可能希望收到通知。让我们添加一个仅在工作流程中的某一步骤发生错误时才会被执行的错误处理机制:
errors:
- id: notify_failure
type: io.kestra.plugin.notifications.slack.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK') }}"
payload: |
{"text": "orders_pipeline在执行过程中失败了,失败编号为{{ execution.id }}"}
现在,当工作流程出现故障时,系统会向Slack频道发送通知,而不会在凌晨3点无声无息地退出执行。需要注意的是,这些敏感信息是通过secret表达式进行保护的。
步骤4:根据事件触发执行,而不仅仅是定时执行
调度任务只是一种触发方式。假设订单并不会按照固定的时间表到达;实际上,只要上游系统决定发送数据,文件就会立即被上传到云存储中。
使用cron调度表来“每隔5分钟检查一次”,这种做法既浪费资源又会导致延迟。而基于事件的触发机制则更为高效:只有在相关事件发生时才执行工作流程。
从概念上来说,我们不需要像前面例子中那样添加定时触发器:
triggers:
- id: every_minute
type: io.kestra.plugin.core.trigger.Schedule
cron: "* * * * *"
相反,我们可以添加一个触发器,该触发器会在S3存储桶中出现新对象时立即执行相应的操作:
triggers:
- id: new_s3_object
type: io.kestra.plugin.aws.s3.Trigger
interval: "PT1M"
accessKeyId: "{{ secret('AWS_ACCESS_KEY_ID') }}"
secretKeyId: "{{ secret('AWS_SECRET_KEY_ID') }}"
region: "eu-central-1"
bucket: "my-bucket"
prefix: "incoming/"
on: CREATE
action: NONE
或者,你也可以设置一个Webhook触发器,通过发送POST请求来启动工作流程的执行;另外,还有实时触发机制,它们可以监听Kafka队列等流式数据源。无论采用哪种触发方式,工作流程的具体内容都是不变的。这样一来,你的自动化系统就能主动响应各种外部事件,而不仅仅被动地等待时间流逝。
步骤5:填补过去的数据
最后,我们来考虑“在处理过程中出现错误”的这种情况。既然你已经修复了计算错误,那么就需要重新运行过去两个月中每一天的数据处理流程。如果使用cron来执行这样的操作,会非常麻烦。
而在Orchestrator中,“补丁操作”是一种针对定时工作流的高级功能:你只需指定开始日期和结束日期,系统就会在该时间段内每隔一定时间间隔自动执行一次任务。每次执行都会通过类似{{ trigger.date }}这样的表达式来识别当前执行的日期,因此你的数据处理步骤可以利用这个日期来获取并处理相应的数据。
在这里,“幂等性”原则就不再只是理论上的概念了。因为补丁操作会重新执行那些可能已经被处理过的任务,所以你的数据写入步骤应该使用基于日期的“插入或替换”机制,这样例如如果3月14日这个任务被执行了两次,数据库中的数据也会和只执行一次时的结果完全相同。
如果你在设计阶段就考虑到这一点,那么补丁操作就会变得常规化,而不会让人感到恐惧。只有当你使用了幂等性强的任务时,第4道屏障(补丁操作问题)才能真正得到解决。
下一步该怎么做
将所有这些内容内化为实际操作的最佳方法,就是把现有的cron作业重新设计成正规的工作流。可以先从只包含一个任务的简单版本开始,确认它能够正常运行并且会在运行记录中显示出来,然后再添加第二个依赖任务、重试机制以及失败警报功能。
每一个步骤都对应着“四道屏障”中的某一条,而当你第一次遇到任务失败、系统自动重试并最终成功恢复时,你就会立刻感受到这种设计的优势——这些机制会确保系统在不出故障的情况下有序地运行。
你并不需要从头开始编写所有的代码。Kestra提供了大量的蓝图模板,这些都是现成的、可以直接复制使用的流程示例,你可以在kestra.io/blueprints网站上查看这些蓝图,或者直接在你的系统中“蓝图”选项卡下找到它们。
每一个蓝图都是一个完整的、可执行的示例,其中会详细说明其功能以及如何对其进行扩展。这样你就可以从接近目标的状态开始进行开发,而不是盲目猜测代码的语法结构。
cron证明了计算机可以在你睡觉的时候继续运行,而Orchestration则更进一步,确保这些计算过程能够可靠地、按顺序地进行,并且其运行结果也一目了然。这种转变意味着,我们的工作重点已经从单纯编写脚本,转变为管理整个系统了。
相关文章
克劳德——注册建筑师培训课程:为参加Anthropic新的认证考试做准备
Anthropic正在扩展其官方认证体系,现在,那些拥有高级技能的从业者可以通过参加“Claude认证架构师——基础课程”考试( CCA-f )来验证自己的专业能力。 不过,阅读官方考试指南可能会让人感到有些困难,因为Anthropic提供的标准文档往往侧重于概念性内容。为了解决这个问题,我们刚刚在freeCodeCamp.org的YouTube频道上发布了一门内容全面、注重实践操作的备考课程,该课程由拥有20年开发经验的Andrew Brown主讲。 这门课程非常重视实际操作环节,而不仅仅是理论讲解。其主要涵盖以下领域: 智能体架构与协调机制: 构建可靠的智能体工作流程,并管理完整的智能体运
阅读全文
演讲主题:将工作流程编译成数据库:这种本应无法实现但却实际可行的架构
Jeremy Edberg与Qian Li探讨了为什么外部协调工具会降低系统的可靠性,以及如何利用现有的数据库来实现可靠的数据处理流程。他们介绍了DBOS Transact是如何利用标准表格、SKIP LOCKED队列以及唯一的主键来管理那些复杂且具备容错能力的人工智能工作流程的——这些机制能够在将延迟降到最低的同时,完全避免额外分布式系统所带来的运营开销。 作者:Jeremy Edberg, Qian Li
阅读全文
通过帮助开发人员来确保平台工程团队遵守相关规范
当一个新的平台团队试图通过缺乏文档支持的、强制性的工作流程来实施他们的规划时,开发人员的体验会明显下降。成功的经验在于简化管理流程、优先处理真正重要的事情,并通过预防、检测和沟通等手段逐步推进各项合规要求的落实。同理心、专注精神以及共同的目标是推动这些措施获得成功实施的关键因素。 作者:本·林德斯
阅读全文