尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Koheesio核心概念速查手册:Step、Reader、Writer与Transformation一网打尽

Koheesio核心概念速查手册:Step、Reader、Writer与Transformation一网打尽 Koheesio核心概念速查手册Step、Reader、Writer与Transformation一网打尽【免费下载链接】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/koheesioKoheesio是一个用于构建高效数据管道的 Python 框架它通过模块化 可复用的设计理念让你能用简单、独立的组件拼装出复杂的数据管道Data Pipeline。无论你是刚接触数据工程的新手还是想快速上手 Koheesio 的普通用户这份核心概念速查手册都能帮你快速理解 Step、Reader、Writer 与 Transformation 四大核心组件少走弯路。一图看懂 Koheesio 的四大核心组件 Koheesio 的整个框架架构建立在四个核心抽象之上它们之间的关系可以用一句话概括Step步骤是所有组件的基类Reader、Writer、Transformation 都是 Step 的特殊形态而 Task任务则是它们的组合。组件一句话理解扮演角色Step可执行的最小逻辑单元一切组件的地基Reader负责读数据数据管道的入口ExtractTransformation负责改数据数据管道的加工车间TransformWriter负责写数据数据管道的出口Load官方文档在 docs/reference/concepts/concepts.md 中用类图清晰展示了这四者的继承与协作关系值得反复研读。Step一切数据管道的积木 Step 到底是什么在 Koheesio 中Step 是一个原子操作是构建数据管道的最小积木。你可以把 Step 想象成一个函数输入Input类的字段field带类型注解输出Output嵌套的Output类继承自StepOutput执行逻辑execute()方法class MyStep(Step): a: str # 输入字段 class Output(StepOutput): # 输出字段 b: str def execute(self) - MyStep.Output: self.output.b f{self.a}-some-suffix这段代码中a是输入b是输出execute是执行逻辑。是不是很像函数但 Step 比普通函数多了一大堆免费福利。用 Step 的 7 大理由新手必看模块化每个 Step 是独立单元出问题能快速定位可复用定义一次随处使用可读性输入/输出/逻辑一目了然自动校验基于 Pydantic输入输出自动类型校验自动日志开始、结束、输入、输出全部自动记录错误处理异常自动捕获、记录、再抛出可扩展天然适配 Apache Spark 等分布式框架源码定义在 src/koheesio/steps/init.py官方详细说明见 docs/reference/concepts/step.md。Reader数据管道的进水口 Reader 的核心职责Reader 是一种特殊的 SparkStep负责从各种数据源读取数据并把结果存到self.output.df一个 Spark DataFrame中。数据源可以是文件、数据库、Web API……几乎任何东西。Reader 的 3 个便捷特性.read()方法等价于.execute().output.df一步拿到 DataFrame.df属性访问输出 DataFrame 的快捷方式如果还没执行会自动触发self.spark随时访问当前活跃的 SparkSessionmy_reader MyReader() df my_reader.read() # 一行搞定直接拿到 DataFrameKoheesio 内置的常用 Reader 速查数据源对应 Reader源码位置CSV / Parquet 等文件FileReadersrc/koheesio/spark/readers/file_loader.pyDelta 表DeltaTableReadersrc/koheesio/spark/readers/delta.pyKafkaKafkaReadersrc/koheesio/spark/readers/kafka.pyJDBC 数据库JdbcReadersrc/koheesio/spark/readers/jdbc.pySnowflakeSnowflakeReadersrc/koheesio/spark/readers/snowflake.pyREST APIRestApiReadersrc/koheesio/spark/readers/rest_api.py内存数据MemoryReadersrc/koheesio/spark/readers/memory.pyReader 的基类定义在 src/koheesio/spark/readers/init.py完整讲解见 docs/reference/spark/readers.md。Transformation数据管道的加工车间 Transformation 的核心职责Transformation 同样是一种 SparkStep它接收一个 DataFrame应用某种变换加列、过滤、聚合……再返回一个 DataFrame。数据在它的手里脱胎换骨。三种 Transformation 类型怎么选类型适用场景特点Transformation对整张表做变换最基础的基类ColumnsTransformation对指定列做变换支持单列/多列循环处理ColumnsTransformationWithTarget变换结果存到新列多了target_column字段比如给指定列加 1用ColumnsTransformation只需写class AddOne(ColumnsTransformation): def execute(self): for column in self.get_columns(): self.output.df self.df.withColumn(column, f.col(column) 1)内置 Transformation 推荐清单 HashUUID5为每行生成 UUID5 哈希src/koheesio/spark/transformations/uuid5.pyDataframeLookup两表关联查询src/koheesio/spark/transformations/lookup.pyRepartition重分区src/koheesio/spark/transformations/repartition.pyRowNumberDedup行号去重src/koheesio/spark/transformations/row_number_dedup.py字符串/日期/重命名等一批现成变换都在 src/koheesio/spark/transformations/ 目录下Transformation 详细文档见 docs/reference/spark/transformations.md。Writer数据管道的出水口 Writer 的核心职责Writer 也是 SparkStep它从self.input.df取出数据写到目标位置文件、数据库、消息队列……。Writer 的便捷特性.write()方法一步完成写入.df属性获取待写入的 DataFrame早期输入校验实例化时就校验输入出错早发现避免白白跑一遍 Spark 任务常用 Writer 一览目标对应 Writer源码位置Delta 表批量DeltaTableWritersrc/koheesio/spark/writers/delta/batch.pyDelta 表SCDSCD2Writersrc/koheesio/spark/writers/delta/scd.pyKafkaKafkaWritersrc/koheesio/spark/writers/kafka.py文件FileWritersrc/koheesio/spark/writers/file_writer.py流式写入StreamWritersrc/koheesio/spark/writers/stream.py注意Writer 还内置了OutputMode枚举如BATCH.APPEND、BATCH.OVERWRITE、BATCH.MERGE定义在 src/koheesio/spark/writers/init.py写 Delta 表时选MERGE模式配合key_columns可实现增量合并。组合拳EtlTask 一网打尽整个管道 ⚡学完四大组件最爽的时刻来了Koheesio 内置了EtlTask把 ReaderExtract TransformationTransform WriterLoad串成一条完整的 ETL 管道代码简洁到令人舒适from koheesio.spark.etl_task import EtlTask etl_task EtlTask( nameMy ETL Task, sourceCsvReader(pathpath/to/source.csv), transformations[Repartition(num_partitions2)], targetDummyWriter(), ) etl_task.execute()这段代码就完成了读取 CSV → 重分区 → 写出 的完整流程。EtlTask的定义在 src/koheesio/spark/etl_task.py它本身也是一个 Step完美体现了用简单组件拼装复杂管道的设计哲学。新手最容易踩的 5 个坑 ⚠️忘了Output嵌套类自定义 Step 时输出必须写在class Output(StepOutput)里否则拿不到结果execute里 return 非 Output 类型execute()总是返回 StepOutput多余返回值会被忽略并告警混淆read()与execute()Reader 用.read()拿 DataFrameWriter 用.write()执行写入Transformation 链式顺序EtlTask.transformations列表顺序就是执行顺序顺序很重要忽略自动校验的好处输入输出类型错误会在早期暴露善用这一点能省大量调试时间总结三分钟记住核心概念 Step 最小积木输入 输出 executeReader 读数据.read()拿 DataFrameTransformation 改数据DataFrame → DataFrameWriter 写数据.write()落地EtlTask 一键组合以上三者想深入学习官方文档的 docs/reference/concepts/step.md 和 docs/reference/concepts/context.md 是必读资料。用 Koheesio 构建数据管道你会发现复杂管道原来可以这么简单【免费下载链接】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),仅供参考
返回列表