公司动态

Koheesio Hello World实战:从零编写你的第一个自定义Step

📅 2026/8/20 16:42:41
Koheesio Hello World实战:从零编写你的第一个自定义Step
Koheesio Hello World实战从零编写你的第一个自定义Step【免费下载链接】koheesioPython framework for building efficient data pipelines. It promotes modularity and collaboration, enabling the creation of complex pipelines from simple, reusable components.项目地址: https://gitcode.com/gh_mirrors/ko/koheesio如果你是数据工程领域的新手想用 Koheesio 快速搭建可复用、易维护的数据管道那么这篇文章就是为你准备的。Koheesio 是一个基于 Python 的数据管道框架它的核心思想是把复杂的 ETL 流程拆解成一个个小而美的自定义 Step步骤。本文将手把手带你完成 Koheesio Hello World 实战从安装环境、理解 Step 结构到编写并运行你的第一个自定义 Step全程零基础可上手代码量极少几分钟就能跑通。Koheesio Step 是什么数据管道的最小积木在 Koheesio 中Step是构建数据管道的最小单元。你可以把它想象成乐高积木一个 Step 接收一组输入经过内部逻辑处理后产生一组输出。管道Pipeline就是由这些积木依次拼接而成的。Koheesio 官方文档 step.md 中给出了清晰的解释Step 是一个原子操作是 Koheesio 构建数据管道的基础构件一个 Task 通常由一系列 Step 组成。Koheesio 内置了多种 Step 类型Reader读取数据、Transformation转换数据、Writer写出数据和Task任务编排它们全部继承自同一个Step基类源码位于 steps/init.py。上图展示了 Koheesio 官方文档系统的分类方式本篇文章属于其中的 Tutorials 实战类内容带你边学边做。环境准备安装 Koheesio 并启动 SparkSession动手之前先做好两件准备工作。Koheesio 要求 Python 3.9 及以上版本入门安装方式可以参考 getting-started.mdpip install koheesio[spark] 我们编写的 HelloWorldStep 基于 Spark所以需要安装带 Spark 扩展的版本。如果只是用纯 Python Step直接pip install koheesio即可。接下来创建 SparkSession。注意Koheesio 不会替你创建 SparkSession每个SparkStep都通过spark属性访问当前会话from pyspark.sql import SparkSession spark SparkSession.builder.getOrCreate()创建好的spark会话可以显式传给 Step也可以不传——Koheesio 会自动尝试获取当前活跃的 SparkSession相关逻辑见 spark/init.py 中的_get_active_spark_session。手写第一个 HelloWorldStep完整代码与运行现在进入正题。我们创建一个名为HelloWorldStep的自定义 Step它的功能很简单接收一条message字符串把它放进一张只有一行数据的 Spark DataFrame 中并输出。完整代码参考 hello-world.mdfrom koheesio.spark import SparkStep class HelloWorldStep(SparkStep): message: str def execute(self) - SparkStep.Output: # 创建一个包含单行数据的 DataFrame self.output.df self.spark.createDataFrame( [(1, self.message)], [id, message] )运行它同样简单hello HelloWorldStep(messageHello, World!) hello.execute() hello.output.df.show()输出结果---------------- | id| message| ---------------- | 1|Hello, World!| ---------------- 你的第一个 Koheesio 自定义 Step 已经跑通了逐步拆解输入、execute 与 Output 三要素上面这段代码虽然短却包含了 Step 的全部核心机制值得逐行理解输入即类属性message: str定义了一个输入字段。Koheesio 基于 Pydantic 构建任何类属性都会在实例化时自动做类型校验。传入非字符串会直接报错帮你尽早发现问题。execute 方法即业务逻辑这是你唯一必须实现的方法里面写这个 Step 要做什么。通过self.message访问输入通过self.output.df写入输出。Output 即结果容器SparkStep自带的Output类中有一个df字段专门用来存放输出的 Spark DataFrame。访问step.output.df就能拿到结果。如果你不想用 Spark也可以直接继承Step写纯 Python Step参考源码 steps/dummy.pyfrom koheesio import Step, StepOutput class DummyStep(Step): a: str b: int class Output(StepOutput): # 自定义输出字段 c: str def execute(self) - None: self.output.c self.a * self.b深入理解为什么 execute 不需要写 return细心的你可能会问execute方法明明没有return为什么还能拿到输出秘密藏在 Koheesio 的元类StepMetaClass中源码见 steps/init.py 的_execute_wrapper。当你定义 Step 类时Koheesio 会自动为execute方法包上一层增强外衣替你完成三件事自动日志运行前后自动打印 Step 的输入与输出方便调试输出校验执行完毕后自动校验输出是否符合Output类定义的类型错误处理异常会被捕获、记录日志后再抛出报错位置一目了然。因此你只需专注于业务逻辑把结果写进self.output剩下的交给框架。进阶技巧让 Step 支持自定义输出字段当内置输出不够用时你可以像DummyStep那样在 Step 内部定义Output类来扩展输出。甚至可以让输入字段拥有默认值和详细描述增强代码可读性from koheesio import Step, StepOutput from koheesio.models import Field class GreetingStep(Step): name: str Field(description要问候的名字) greeting: str Field(defaultHello, description问候语) class Output(StepOutput): full_greeting: str def execute(self) - None: self.output.full_greeting f{self.greeting}, {self.name}!调用时输出会自动校验并记录日志一个带类型安全、自动日志的组件就完成了。从 Step 到管道如何组合出完整 ETL 流程单个 Step 只是起点。Koheesio 的强大之处在于组合你可以把自定义 Step 与内置的 Reader、Transformation、Writer 拼成一条完整管道。Readerkoheesio.spark.readers下的内置读取器如DummyReader、DeltaReaderTransformationkoheesio.spark.transformations下的转换步骤如CamelToSnakeTransformationWriterkoheesio.spark.writers下的写出器如DummyWriterTaskEtlTasketl_task.py帮你把读取 → 转换 → 写出编排成完整任务。在 hello-world.md 中还有一个MyFavoriteMovieTask的完整示例展示了如何用EtlTask组合 Reader 与 Writer配合 sample.yaml 配置文件构建真实管道。想系统学习 Koheesio 的核心概念可以阅读 learn-koheesio.md。常见问题排查Q1报错 No active SparkSession foundKoheesio 不会自动创建 SparkSession。请先运行SparkSession.builder.getOrCreate()或把spark作为参数传入 Step。Q2输入传错类型怎么办Pydantic 会自动校验并抛出清晰的类型错误检查传入参数的类型即可。Q3execute里可以写return吗可以但不必要。框架会自动包装返回值你只需把结果写入self.output。总结与更多学习资源至此你已经完成了 Koheesio Hello World 实战掌握了编写自定义 Step 的核心三要素输入字段、execute 方法、Output 输出。下一步可以挑战更复杂的转换逻辑或尝试用EtlTask组装完整的数据管道。想继续深入推荐阅读项目内这些资料入门教程getting-started.mdStep 概念详解step.md完整示例hello-world.md学习路线learn-koheesio.md现在动手写属于你的第一个自定义 Step 吧【免费下载链接】koheesioPython framework for building efficient data pipelines. It promotes modularity and collaboration, enabling the creation of complex pipelines from simple, reusable components.项目地址: https://gitcode.com/gh_mirrors/ko/koheesio创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考