ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Feast SparkSource 数据源深度指南:从表/查询/文件引用到 Iceberg、Delta Lake 与 Hudi 数据湖表格式

Feast SparkSource 数据源深度指南:从表/查询/文件引用到 Iceberg、Delta Lake 与 Hudi 数据湖表格式 Feast SparkSource 数据源深度指南从表/查询/文件引用到 Iceberg、Delta Lake 与 Hudi 数据湖表格式【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast本篇指南完整讲解 Feast 中基于 Spark 的离线数据源SparkSource如何通过表引用、SQL 查询与文件路径三种方式接入数据如何配套SparkOfflineStoreConfig配置 Spark 会话以及如何借助IcebergFormat、DeltaFormat、HudiFormat三种表格式将数据湖表Data Lakehouse直接接入特征存储。读完本文你将掌握SparkSource的全部构造参数、参数组合约束、底层读取逻辑与类型映射规则并能在自己的 feature repo 中直接写出可运行的 Spark 数据源定义。一、SparkSource 是什么SparkSource是 Feast 中由社区贡献contrib实现的 Spark 离线数据源定义于 sdk/python/feast/infra/offline_stores/contrib/spark_offline_store/spark_source.py。它的数据来源可以是某个 Spark 存储如 Hive Metastore 或内存中注册的临时表中的表一段可在 Spark SQL 中执行的查询一组位于本地或云存储S3、HDFS 等上的文件。从 Feast 的角度看SparkSource是DataSource的子类通过source_type()返回DataSourceProto.BATCH_SPARK即它是一种批式数据源必须与 Spark offline store 配合使用。在 docs/reference/data-sources/overview.md 的功能矩阵中Spark 支持八种基本类型、数组类型、Map与Struct类型详见后文类型映射一节。状态与免责声明需要特别提醒的是官方文档明确给出两条 DisclaimerSpark 数据源尚未达到完整的测试覆盖率不要假设其完全稳定源码中也对使用者发出RuntimeWarning提示该 API 处于 alpha 开发阶段、未来可能变化见 spark_source.py 与 spark.py。因此它适合评估、试点与有相应兜底机制的生产场景而不建议在没有充分测试的前提下作为核心依赖。二、三种基础引用方式SparkSource支持通过table、query、path三种方式定位数据三种方式互斥且有严格校验校验逻辑见 SparkOptions。1. 表引用table适用于数据已注册在 SparkSession 中内存临时表或 Hive Metastore 表的场景from feast.infra.offline_stores.contrib.spark_offline_store.spark_source import ( SparkSource, ) my_spark_source SparkSource( tableFEATURE_TABLE, )底层在生成 SQL 时会为表名中的每个部分加上反引号以安全引用..join([f{x} for x in self.table.split(.)])因此多级命名空间如db.feature_table也可以直接传入。2. SQL 查询query适用于需要先经过加工、过滤或聚合再作为特征源的场景from feast.infra.offline_stores.contrib.spark_offline_store.spark_source import ( SparkSource, ) my_spark_source SparkSource( querySELECT timestamp as ts, created, f1, f2 FROM spark_table, )在get_table_query_string()中查询会被包裹成子查询({self.query})嵌入后续 SQL。3. 文件引用path适用于读取 Parquet、Avro、CSV、JSON 等文件。此时必须配合file_format指定文件格式from feast.infra.offline_stores.contrib.spark_offline_store.spark_source import ( SparkSource, ) my_spark_source SparkSource( pathf{CURRENT_DIR}/data/driver_hourly_stats, file_formatparquet, timestamp_fieldevent_timestamp, created_timestamp_columncreated, )读取文件时Feast 会先把文件加载为 DataFrame注册为临时视图df.createOrReplaceTempView(...)再以该视图作为后续 SQL 的数据源。若文件路径无法读取会抛出DataSourceNotFoundException。参数组合约束源码级SparkOptions.__init__中实现的校验规则非常明确写配置时务必遵守组合结果同时指定table与query/path抛出ValueError: table cannot be combined with query or pathtable、query、path三者都为空抛出ValueError: At least one of params(table, query, path) must be specified指定path但既无table_format又无file_format抛出ValueErrorfile_format必填指定file_format但不在{csv, json, parquet, avro}之列抛出ValueError见SparkFileSourceFormat枚举其中query path是允许的组合query 用于物化materialization期间的读取path 用于离线写回offline_write_batch与get_historical_features。三、完整参数说明SparkSource.__init__支持的全部参数见 spark_source.py参数类型说明namestr数据源名称项目内唯一未指定时默认取table的值。若name与table均为空抛出DataSourceNoNameExceptiontablestrSpark 表名支持db.table形式querystr在 Spark 中执行的 SQL 查询pathstr文件数据路径file_formatstr底层文件格式parquet、avro、csv、jsontable_formatTableFormat表元数据格式iceberg / delta / hudi与file_format相互独立、可选timestamp_fieldstr事件时间戳字段用于特征值的 point-in-time joincreated_timestamp_columnstr记录创建时间戳字段用于行去重同一实体多条记录时取最新date_partition_columnstr日期分区列用于按分区裁剪、加速大数据量下的检索date_partition_column_formatstr分区列的时间格式默认%Y-%m-%dfield_mappingdict数据源列名到特征名称的映射descriptionstr人类可读描述tagsdict任意键值对元数据ownerstr数据源负责人通常是维护者邮箱connection_refConnectionRef数据源级别的外部凭据引用用于多租户场景详见 docs/reference/data-sources/overview.md这些参数会序列化进DataSourceProto.SparkOptions在feast apply时写入注册表并在读取时通过from_proto还原。四、配套的 Spark 离线存储配置SparkSource必须搭配SparkOfflineStore使用。在 feature repo 的feature_store.yaml中offline_store一节支持以下配置项见 spark.pyproject: my_project registry: data/registry.db provider: local offline_store: type: spark spark_conf: spark.sql.session.timeZone: UTC spark.sql.parser.quotedRegexColumnNames: true spark.ui.enabled: false staging_location: s3://my-bucket/staging region: eu-west-1 online_store: type: sqlite配置项说明type固定为sparkspark_confSpark 会话的配置覆盖项以键值对形式传给SparkConf().setAll(...)staging_location批式物化materialization作业的远程暂存路径region适用于 S3 类暂存位置时的 AWS 区域会话获取逻辑在 get_spark_session_or_start_new_with_repoconfig 中优先复用SparkSession.getActiveSession()否则用spark_conf构建新会话并强制设置spark.sql.parser.quotedRegexColumnNamestrue保证带反引号的列名与正则列名能被正确解析。另外SparkOfflineStore.supports_filter_by_created_timestamp True即支持按created_timestamp_column过滤它还实现了compute_monitoring_metrics、get_monitoring_max_timestamp、ensure_monitoring_tables等监控能力可为特征做数据质量统计。五、数据湖表格式支持Iceberg / Delta Lake / Hudi新特性SparkSource现已支持 Apache Iceberg、Delta Lake 与 Apache Hudi 三种高级表格式带来 ACID 事务、时间旅行time travel与 schema 演化能力。三种格式的配置类统一继承自 sdk/python/feast/table_format.py 中的抽象基类TableFormatTableFormatType枚举iceberg/delta/hudi可通过set_property(key, value)动态增改属性也有to_proto/from_proto/to_dict/from_dict/from_json等序列化方法。1. Apache IcebergIceberg 专为海量分析数据集设计提供 ACID 事务快照隔离的原子提交、时间旅行、安全的 schema 演化与隐藏分区。from feast.infra.offline_stores.contrib.spark_offline_store.spark_source import SparkSource from feast.table_format import IcebergFormat iceberg_format IcebergFormat( catalogmy_catalog, namespacemy_database ) my_spark_source SparkSource( nameuser_features, pathmy_catalog.my_database.user_table, table_formaticeberg_format, timestamp_fieldevent_timestamp )配置参数参数类型说明catalogstr可选Iceberg catalog 名称会写入属性iceberg.catalognamespacestr可选catalog 内的命名空间/库会写入属性iceberg.namespacepropertiesdict可选额外的 Iceberg 配置属性常用属性iceberg_format IcebergFormat( catalogspark_catalog, namespaceproduction, properties{ # 快照选择 snapshot-id: 123456789, as-of-timestamp: 1609459200000, # Unix 毫秒时间戳 # 性能调优 read.split.target-size: 134217728, # 128 MB 分片 read.parquet.vectorization.enabled: true, # 高级配置 io-impl: org.apache.iceberg.hadoop.HadoopFileIO, warehouse: s3://my-bucket/warehouse } )时间旅行示例# 读取特定快照 iceberg_format IcebergFormat(catalogspark_catalog, namespacelakehouse) iceberg_format.set_property(snapshot-id, 7896524153287651133) # 或按时间点读取 iceberg_format.set_property(as-of-timestamp, 1609459200000)2. Delta LakeDelta Lake 是开源的存储层为 Spark 与大数据负载带来 ACID 事务、时间旅行、schema 强制与统一批流处理。from feast.infra.offline_stores.contrib.spark_offline_store.spark_source import SparkSource from feast.table_format import DeltaFormat delta_format DeltaFormat() my_spark_source SparkSource( nametransaction_features, paths3://my-bucket/delta-tables/transactions, table_formatdelta_format, timestamp_fieldtransaction_timestamp )配置参数参数类型说明checkpoint_locationstr可选Delta 事务日志 checkpoint 位置会写入属性delta.checkpointLocation流式场景必需propertiesdict可选额外的 Delta 配置属性常用属性delta_format DeltaFormat( checkpoint_locations3://my-bucket/checkpoints, properties{ # 时间旅行 versionAsOf: 5, timestampAsOf: 2024-01-01 00:00:00, # 性能优化 delta.autoOptimize.optimizeWrite: true, delta.autoOptimize.autoCompact: true, # 数据跳过 delta.dataSkippingNumIndexedCols: 32, # Z-ordering delta.autoOptimize.zOrderCols: event_timestamp } )时间旅行示例# 按版本读取 delta_format DeltaFormat() delta_format.set_property(versionAsOf, 10) # 或按时间戳读取 delta_format DeltaFormat() delta_format.set_property(timestampAsOf, 2024-01-15 12:00:00)3. Apache HudiHudiHadoop Upserts Deletes and Incrementals擅长增量数据处理支持高效记录级更新、增量查询与时间旅行。from feast.infra.offline_stores.contrib.spark_offline_store.spark_source import SparkSource from feast.table_format import HudiFormat hudi_format HudiFormat( table_typeCOPY_ON_WRITE, record_keyuser_id, precombine_fieldupdated_at ) my_spark_source SparkSource( nameuser_profiles, paths3://my-bucket/hudi-tables/user_profiles, table_formathudi_format, timestamp_fieldevent_timestamp )配置参数参数类型说明table_typestr可选COPY_ON_WRITE或MERGE_ON_READ写入属性hoodie.datasource.write.table.typerecord_keystr可选唯一标识记录的一个或多个字段复合键用逗号分隔写入hoodie.datasource.write.recordkey.fieldprecombine_fieldstr可选决定记录最新版本的字段通常是时间戳/版本号写入hoodie.datasource.write.precombine.fieldpropertiesdict可选额外的 Hudi 配置属性两种表类型的选择COPY_ON_WRITECOW数据以列式Parquet存储更新生成新文件版本适合读多写少、查询延迟低的场景MERGE_ON_READMOR列式 行式混合存储更新先写入 delta 日志适合写多读少、写入延迟低的场景。常用属性与增量查询hudi_format HudiFormat( table_typeCOPY_ON_WRITE, record_keyuser_id, precombine_fieldupdated_at, properties{ # 查询类型 hoodie.datasource.query.type: snapshot, # 或 incremental # 增量查询区间 hoodie.datasource.read.begin.instanttime: 20240101000000, hoodie.datasource.read.end.instanttime: 20240102000000, # 索引 hoodie.index.type: BLOOM, # MOR 表压缩 hoodie.compact.inline: true, hoodie.compact.inline.max.delta.commits: 5, # 聚类 hoodie.clustering.inline: true } )# 只处理新增/变化的数据 hudi_format HudiFormat( table_typeCOPY_ON_WRITE, record_keyid, precombine_fieldtimestamp, properties{ hoodie.datasource.query.type: incremental, hoodie.datasource.read.begin.instanttime: 20240101000000, hoodie.datasource.read.end.instanttime: 20240102000000 } )4. 文件格式与表格式两者可叠加使用需要区分两个层次文件格式是数据的物理编码Parquet、Avro、CSV、JSON表格式是构建其上的元数据与事务层Iceberg、Delta、Hudi。二者可以同时使用from feast.infra.offline_stores.contrib.spark_offline_store.spark_source import SparkSource from feast.table_format import IcebergFormat iceberg IcebergFormat(catalogmy_catalog, namespacedb) source SparkSource( namefeatures, pathcatalog.db.table, file_formatparquet, # 底层存储格式 table_formaticeberg, # 表元数据格式 timestamp_fieldevent_timestamp )从源码看读取路径在 spark_source.py 的_load_dataframe_from_path中分派未指定table_format时走spark_session.read.format(self.file_format).load(self.path)指定后则用table_format.format_type.value即iceberg/delta/hudi作为 Spark read format并把table_format.properties逐个作为 reader option 传入。5. 表格式选择建议使用场景推荐格式理由频繁 schema 变更的大规模分析Icebergschema 演化与隐藏分区能力强生态成熟流批一体的负载Delta Lake与 Spark 集成紧密统一批流CDC 与大量 upsertHudi高效记录级更新与增量查询读多分析Iceberg 或 Delta查询性能出色写多事务型HudiMOR写入路径优化多引擎互通Iceberg引擎支持面最广Spark、Flink、Trino 等详细配置说明与最佳实践分区策略、自动优化、历史清理、元数据监控等可继续阅读 docs/reference/data-sources/table-formats.md。六、类型映射Spark 类型与 Feast 类型SparkSource.source_datatype_to_feast_value_type()返回spark_to_feast_value_type该函数定义于 sdk/python/feast/type_map.py实现 Spark 类型字符串到 FeastValueType的映射Spark 类型Feast 类型stringSTRINGint/shortINT32bigint/longINT64floatFLOATdouble/decimalDOUBLEbooleanBOOLtimestamp/dateUNIX_TIMESTAMPbyteBYTESarraybyte等 8 类数组对应的*_LIST类型map...MAParraymap...MAP_LISTstruct...STRUCTarraystruct...STRUCT_LIST这也解释了官方文档中支持全部八种基本类型及其数组类型的说法并进一步覆盖了Map与Struct。未识别的类型如interval、binary会回退为NULL。此外spark_schema_to_np_dtypestype_map.py还负责把 Spark schema 转换为 NumPy dtype用于历史特征检索结果落成 DataFrame。完整的功能对比矩阵见 docs/reference/data-sources/overview.md。七、与 Feature View 的组合使用SparkSource通常作为FeatureView或BatchFeatureView的batch_source。示例参考测试仓库 tests/data_source.py 中SparkDataSourceCreator.create_data_source的构造方式from datetime import timedelta from feast import FeatureView, Field from feast.types import Float32, Int64 from feast.infra.offline_stores.contrib.spark_offline_store.spark_source import SparkSource driver_stats_source SparkSource( namedriver_stats, pathdata/driver_hourly_stats, file_formatparquet, timestamp_fieldevent_timestamp, created_timestamp_columncreated, ) driver_stats_fv FeatureView( namedriver_stats, entities[driver_id], ttltimedelta(days7), schema[ Field(nameconv_rate, dtypeFloat32), Field(nameacc_rate, dtypeFloat32), Field(namedriver_id, dtypeInt64), ], onlineTrue, batch_sourcedriver_stats_source, )在历史特征检索get_historical_features时SparkOfflineStore会上传实体 DataFrame 为临时视图 → 按特征视图生成 point-in-time join SQL见build_point_in_time_query→ 返回SparkRetrievalJob执行并落成 DataFrame。物化materialization则通过pull_latest_from_table_or_query用ROW_NUMBER() OVER (PARTITION BY join_key ORDER BY timestamp DESC)取每个实体的最新特征行并支持date_partition_column按分区裁剪过滤。八、测试与验证Spark 数据源的通用测试位于 sdk/python/feast/infra/offline_stores/contrib/spark_offline_store/tests/data_source.py其SparkDataSourceCreator使用本地模式masterlocal[*]启动 SparkSession并设置了 UTC 时区、关闭 UI/事件日志等测试友好的配置通过写 Parquet 临时文件构造SparkSource来跑通特征检索与物化流程。注册配置见 spark_repo_configuration.py。这也是评估本数据源稳定性最直接的方式——在你自己的环境复现这些测试再决定是否投入生产。九、延伸阅读docs/reference/data-sources/table-formats.mdIceberg / Delta / Hudi 的完整属性清单与最佳实践docs/reference/data-sources/overview.md各批式数据源功能矩阵与ConnectionRef凭据机制sdk/python/feast/infra/offline_stores/contrib/spark_offline_store/spark.pySpark 离线存储的完整实现物化、历史检索、监控sdk/python/feast/infra/offline_stores/contrib/spark_offline_store/spark_source.pySparkSource与SparkOptions的完整源码sdk/python/feast/table_format.py三种表格式配置类的实现与工厂函数【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表