完整使用指南)
Feast 中 ClickHouse 数据源ClickhouseSource完整使用指南【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast导读本文围绕 Feast 开源特征存储The Open Source Feature Store for AI/ML中的 ClickHouse 批处理数据源展开系统讲解如何通过ClickhouseSource将 ClickHouse 的表或视图接入 Feast 特征仓库涵盖定义方式、参数说明、类型映射、离线存储搭配与底层实现原理。读完本文你将能够独立完成 ClickHouse 数据源的定义、feature_store.yaml离线存储配置并理解点即时point-in-time连接、最新值拉取等核心流程在 ClickHouse 上是如何落地的。一、ClickHouse 数据源概述ClickhouseSource是 Feast 社区贡献contrib的批处理数据源实现代码位于 clickhouse_source.py。从数据模型上看ClickHouse 数据源的本质是 ClickHouse 中的表table或视图view。它有两种指定方式二者选其一表引用table reference直接给出 ClickHouse 中的表名或视图名SQL 查询SQL query给出任意合法的 SELECT 查询Feast 会将其作为子查询包装使用。这种“表或查询”的双通道设计使数据源既可以直接指向已物化的宽表也可以通过 SQL 灵活裁剪列、过滤行或完成预聚合而不必在 ClickHouse 中额外建表。需要特别注意的是该数据源属于社区贡献组件与FileSource、BigQuerySource等核心数据源相比尚未达到完整的测试覆盖率官方不保证其完全稳定在生产环境中使用前应自行做好充分验证。二、安装前提要在 Feast 中使用 ClickHouse 数据源以及配套的 ClickHouse 离线存储需要安装带有 ClickHouse 扩展依赖的 Feastpip install feast[clickhouse]该扩展会引入clickhouse-connect等连接 ClickHouse 所需的依赖包对应实现见 clickhouse.py 与 connection_utils.py。三、定义一个 ClickHouse 数据源官方文档给出的最小定义示例如下from feast.infra.offline_stores.contrib.clickhouse_offline_store.clickhouse_source import ( ClickhouseSource, ) driver_stats_source ClickhouseSource( namefeast_driver_hourly_stats, querySELECT * FROM feast_driver_hourly_stats, timestamp_fieldevent_timestamp, created_timestamp_columncreated, )这里使用query参数把整张表包装成查询子句。若希望直接引用表名可以改用table参数driver_stats_source ClickhouseSource( namefeast_driver_hourly_stats, tablefeast_driver_hourly_stats, timestamp_fieldevent_timestamp, created_timestamp_columncreated, )3.1 构造参数说明对照源码 clickhouse_source.py 中的__init__签名ClickhouseSource支持以下参数参数类型默认值说明nameOptional[str]None数据源名称用于在 Feast 注册表中标识该数据源若未提供则回退使用table值queryOptional[str]NoneClickHouse SQL 查询语句与table二选一tableOptional[str]NoneClickHouse 表或视图名与query二选一timestamp_fieldOptional[str]事件时间戳字段名用于点即时连接与 TTL 裁剪created_timestamp_columnOptional[str]创建时间戳列用于同事件时间下多条记录的去重取最新field_mappingOptional[dict[str, str]]None原始列名到特征列名的映射descriptionOptional[str]数据源描述tagsOptional[dict[str, str]]None附加标签如关联的ConnectionRef凭据信息ownerOptional[str]数据源负责人connection_refOptional[ConnectionRef]None独立于全局配置的外部凭据引用用于多租户/凭据隔离场景3.2 命名与校验规则源码中有一条强制校验逻辑当name与table同时为空时直接抛出DataSourceNoNameException见 clickhouse_source.pyif name is None and table is None: raise DataSourceNoNameException() name name or table也就是说name与table至少要提供其中一个当仅提供table时name自动取表名。反过来仅提供query时必须显式给出name因为此时没有表名可以回退。3.3 表引用与查询子句的生成ClickhouseSource内部把name/query/table打包进ClickhouseOptions对象同样定义于 clickhouse_source.py并在生成 SQL 时通过get_table_query_string()决定使用哪种形态def get_table_query_string(self) - str: if self._clickhouse_options._table: return f{self._clickhouse_options._table} else: return f({self._clickhouse_options._query})可以看到指定了table就直接使用表名否则把query用括号包裹成子查询。这个字符串在后续get_historical_features、pull_latest_from_table_or_query、pull_all_from_table_or_query等流程中被反复引用是数据源与离线存储协作的关键接口。3.4 序列化与注册表存储ClickhouseSource以自定义数据源CUSTOM_SOURCE的形式写入 Feast 注册表_to_proto_impl将配置序列化为DataSourceProto并记录data_source_class_type为完整类路径见 clickhouse_source.pyfrom_proto负责从注册表反序列化恢复。这意味着定义好数据源后执行feast apply即可把配置持久化到注册表供训练与在线服务阶段复用。四、支持的数据类型与映射规则Feast 内部有一套类型系统包含八种原始类型bytes、string、int32、int64、float32、float64、bool、timestamp以及对应的数组Array类型详见 类型系统文档。4.1 官方支持范围根据官方说明ClickHouse 数据源支持全部八种原始类型及其对应的数组类型ClickHouse 的Decimal类型不直接映射为 Feast 的 Decimal而是通过转换为 double 实现支持。4.2 源码级类型映射类型映射函数ch_type_to_feast_value_type位于 clickhouse_source.py它把 ClickHouse 类型解析后映射到 Feast 的ValueType。映射关系如下ClickHouse 类型Feast ValueTypeBooleanBOOLStringSTRINGFloat32FLOATFloat64DOUBLEDecimalDOUBLE经转换Int32INT32Int64INT64DateTimeUNIX_TIMESTAMPDateTime64UNIX_TIMESTAMPArray(Boolean)BOOL_LISTArray(String)STRING_LISTArray(Float32)FLOAT_LISTArray(Float64)DOUBLE_LISTArray(Decimal)DOUBLE_LISTArray(Int32)INT32_LISTArray(Int64)INT64_LISTArray(DateTime)UNIX_TIMESTAMP_LISTArray(DateTime64)UNIX_TIMESTAMP_LIST从源码可以看到两个值得注意的实现细节Decimal 的双重转换不仅在类型推断阶段把Decimal映射为DOUBLE在ClickhouseRetrievalJob._to_arrow_internal中还会对 Arrow 结果里所有decimal类型的列执行pa.cast(..., target_typepa.float64())强制转换为 double见 clickhouse.py源码注释明确说明“Feast doesnt support native decimal types”。未知类型的处理若映射表中不存在对应关系会返回ValueType.UNKNOWN并打印unknown type:提示这要求你在使用非常规 ClickHouse 类型如UUID、LowCardinality包装类型等时特别谨慎建议先小规模验证。4.3 与其他批数据源的类型支持对比不同批处理数据源对复杂类型的支持存在差异完整对比矩阵见>project: my_project registry: data/registry.db provider: local offline_store: type: feast.infra.offline_stores.contrib.clickhouse_offline_store.clickhouse.ClickhouseOfflineStore host: DB_HOST port: DB_PORT database: DB_NAME user: DB_USERNAME password: DB_PASSWORD use_temporary_tables_for_entity_df: true online_store: path: data/online_store.db对应的离线存储文档见 ClickHouse 离线存储。5.1 离线存储配置参数offline_store段解析为ClickhouseOfflineStoreConfig它继承自 ClickhouseConfig可用参数如下参数类型默认值说明typeLiteral[clickhouse]clickhouse离线存储类型标识hoststr必填ClickHouse 服务器地址portint8123ClickHouse HTTP 端口databasestr必填默认数据库名userstr必填用户名passwordstr必填密码use_temporary_tables_for_entity_dfboolTrue实体 DataFrame 是否以临时表方式上传additional_client_argsdictNone透传给clickhouse_connect.get_client的额外参数如send_receive_timeout、compress等其中use_temporary_tables_for_entity_df为可选参数默认开启。当其为True时实体 DataFrame 通过CREATE TEMPORARY TABLE创建临时表为False时则创建ENGINE MergeTree()的实体表并在任务结束后DROP TABLE见 clickhouse.py。在多线程并发执行大量历史特征任务时可根据 ClickHouse 集群的临时表能力权衡该参数。5.2 实体 DataFrame 的两种形态该离线存储对实体entity数据支持两种输入方式官方文档说明见 ClickHouse 离线存储SQL 查询字符串直接把实体查询作为子查询使用Pandas DataFrame默认上传为 ClickHouse 临时表后再参与连接运算。在源码_upload_entity_df中SQL 字符串形态会执行CREATE [TEMPORARY] TABLE ... AS (entity_df)DataFrame 形态则先依据 Arrow schema 建表_df_to_create_table_schema再insert_df写入两者最终都统一为表引用供点即时连接使用见 clickhouse.py。六、连接管理与线程安全Feast 通过 connection_utils.py 中的get_client获取 ClickHouse 客户端其核心设计是基于threading.local()的线程本地缓存thread_local threading.local() def get_client(config: ClickhouseConfig) - Client: if not hasattr(thread_local, clickhouse_client): ... thread_local.clickhouse_client clickhouse_connect.get_client(...) return thread_local.clickhouse_client由于 clickhouse-connect 的客户端不是线程安全的同一个线程内复用同一实例不同线程各自持有独立实例从而避免并发访问导致崩溃。这一点在单元测试 test_clickhouse.py 中有专门验证同一线程两次调用返回同一对象不同线程返回不同对象。此外ClickhouseConfig.additional_client_args会被展开为clickhouse_connect.get_client的关键字参数如send_receive_timeout、compress、client_name其解析行为也有对应测试覆盖见 test_clickhouse.py。七、核心工作流与底层原理7.1 历史特征获取点即时正确性连接ClickhouseOfflineStore.get_historical_features实现了点即时point-in-time正确的特征连接核心逻辑是解析实体数据DataFrame 或 SQL推断实体时间戳列并计算时间范围将实体数据上传为 ClickHouse 表复用 PostgreSQL 离线存储的build_point_in_time_query构造器配合 ClickHouse 专用的MULTIPLE_FEATURE_VIEW_POINT_IN_TIME_JOINSQL 模板见 clickhouse.py生成多特征视图的连接查询返回惰性执行的ClickhouseRetrievalJob。该 SQL 模板通过 CTE 链完成构造实体entity_row_unique_id→ 按 TTL 与时间窗口过滤特征子查询 → 按event_timestamp排序取ROW_NUMBER() 1的最新记录 → 多特征视图LEFT JOIN合并。模板中还包含对created_timestamp_column的专门去重 CTE__dedup保证同一事件时间下优先取创建时间最新的记录。值得注意的是ClickhouseOfflineStore声明了supports_filter_by_created_timestamp True即支持按创建时间戳过滤历史特征对应测试见 test_filter_by_created_timestamp.py。7.2 最新特征拉取物化到在线存储pull_latest_from_table_or_query用于将最新特征值拉取并物化到在线存储其生成的关键 SQL 模式为SELECT {feature_fields} FROM ( SELECT {fields}, ROW_NUMBER() OVER(PARTITION BY {join_keys} ORDER BY {timestamp_field} DESC, {created_timestamp_column} DESC) AS _feast_row FROM {table_or_query} a WHERE a.{timestamp_field} BETWEEN toDateTime64({start}, 6, {tz}) AND toDateTime64({end}, 6, {tz}) ) b WHERE _feast_row 1即按实体键分区、按时间戳倒序取每个分区的第一条记录。时间过滤使用toDateTime64保留微秒精度见 clickhouse.py。7.3 结果导出与 Arrow 类型转换ClickhouseRetrievalJob继承自PostgreSQLRetrievalJob支持导出为 DataFramequery_df、Arrow 表query_arrow与 SQL并实现了persist方法把结果写回 ClickHouse 表配合SavedDatasetClickhouseStorage见 clickhouse_source.py。在_to_arrow_internal中结果中的 decimal 列会被统一转换为float64与第四节所述 Decimal → double 的类型策略保持一致。7.4 非实体检索模式get_historical_features还支持entity_dfNone的非实体检索模式此时通过compute_non_entity_date_range计算特征视图的时间范围并用结束时间构造单行实体 DataFrame 完成全表历史检索见 clickhouse.py。该路径有专门的单元测试覆盖见 test_clickhouse.py。八、功能能力一览ClickhouseOfflineStore对离线存储标准接口的支持情况官方文档矩阵见 ClickHouse 离线存储接口方法ClickHouse 支持get_historical_features点即时正确连接✅pull_latest_from_table_or_query拉取最新特征值✅pull_all_from_table_or_query拉取已保存数据集❌offline_write_batch批量写入离线存储❌write_logged_features持久化特征日志❌ClickhouseRetrievalJob的能力矩阵检索任务能力ClickHouse 支持导出 DataFrame✅导出 Arrow 表✅导出 Arrow 批次❌导出 SQL✅导出数据湖S3、GCS 等✅导出数据仓库✅导出 Spark DataFrame❌本地执行 Python 按需变换✅远程执行 Python 按需变换❌结果持久化到离线存储✅执行前预览查询计划✅读取分区数据✅可以看到点即时连接与最新值拉取是 ClickHouse 离线存储的核心能力足以支撑训练集生成与在线存储物化的主链路而数据集拉取、批量写入与特征日志持久化暂不支持涉及这些场景时需考虑其他离线存储方案。九、使用建议与注意事项明确 contrib 定位ClickHouse 数据源与离线存储均未达到完整测试覆盖官方不保证完全稳定。从源码结构看其点即时连接模板复用了 PostgreSQL 离线存储的构造器build_point_in_time_query与PostgreSQLRetrievalJob依赖关系较深升级 Feast 版本时需关注兼容性。类型选择优先使用上表列出的九种基础类型及其数组形式Decimal会被强制转为 double涉及精度敏感的金额等字段时应提前确认精度损失是否可接受。命名约束name与table至少提供一个使用query定义数据源时必须显式命名。线程安全Feast 为每个线程维护独立的 ClickHouse 客户端实例请勿在业务代码中跨线程复用通过get_client获取的对象。实体表策略默认使用临时表承载实体数据若 ClickHouse 集群限制临时表使用可将use_temporary_tables_for_entity_df设为false任务结束会自动清理实体表。相关资源数据源实现clickhouse_source.py离线存储实现clickhouse.py连接配置clickhouse_config.py 与 connection_utils.py单元测试test_clickhouse.py配套文档ClickHouse 离线存储、数据源总览、类型系统【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考