重要
此功能目前以公共预览版提供。 工作区管理员可以从 预览 页控制对此功能的访问。 请参阅 Manage Azure Databricks 预览版。
流是一种外部流式数据源,例如 Apache Kafka。 流存储连接详细信息、身份验证、架构和引入配置。 创建流后,可以使用 功能视图 定义来引用它,以创建实时流式处理功能。
流具有三部分名称(catalog.schema.stream_name)。 对流的访问权限由其关联的引入表决定。 有关详细信息,请参阅 引入和回填 。
要求
- 对于运行笔记本命令:无服务器或运行 Databricks Runtime 17.0 ML 或更高版本的经典计算群集。
-
feature-engineering-client必须安装Python包版本 0.16.0 或更高版本。
创建流
用于 create_stream() 创建新的流。 流需要四个配置组件:
- 源配置:指定流式平台(例如 Kafka)以及源特有的详细信息(例如 Kafka 的主题订阅信息)。
- 连接配置:指定如何连接到流平台并进行身份验证,包括引导服务器和凭据。
- 架构配置:定义消息键和值的结构。
- 引入配置:指定流数据引入的位置和方式。 有关详细信息,请参阅 引入和回填 。
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
DirectSchemas,
SchemaConfig,
IngestionConfig,
IngestionDestination,
StreamBackfillSource,
)
client = FeatureEngineeringClient()
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="events-topic"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "transaction_id": {"type": "string"},'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string", "format": "date-time"}'
' }'
'}'
)
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
),
)
连接到流源
在定义流式处理功能之前,请连接并测试与 Kafka 中转站的流式处理 Lakeflow 管道连接。 请参阅 无服务器计算上的流处理 和 连接到 Apache Kafka。
有关 AWS 托管流服务(Amazon MSK),请参阅 到 Amazon MSK 的无服务器私有连接。 有关 Kafka 身份验证选项的详细信息,请参阅 “身份验证”。
Authentication
Unity Catalog 连接(推荐)
使用 Unity 目录连接对 Kafka 群集进行身份验证。 这是用于托管身份验证的推荐方法。 若要创建连接,请参阅 “创建连接”。
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
直接 mTLS
对于直接 mTLS 身份验证,请提供存储在 Unity Catalog 卷上的密钥存储和信任存储文件,并通过 Databricks 机密范围引用密码。 有关使用 Kafka 进行 SSL 身份验证的详细信息,请参阅使用 SSL 将Azure Databricks连接到 Kafka。
from databricks.feature_engineering.entities import (
DirectMtlsConfig,
MtlsConfig,
SecretScopeReference,
)
connection_config = DirectMtlsConfig(
bootstrap_servers="broker1:9092,broker2:9092",
mtls_config=MtlsConfig(
keystore_location="/Volumes/my_catalog/my_schema/my_volume/keystore.jks",
keystore_password_ref=SecretScopeReference(
scope="my_scope", key="keystore_password"
),
key_password_ref=SecretScopeReference(
scope="my_scope", key="key_password"
),
truststore_location="/Volumes/my_catalog/my_schema/my_volume/truststore.jks",
truststore_password_ref=SecretScopeReference(
scope="my_scope", key="truststore_password"
),
),
)
SASL
预览版期间不支持 SASL 身份验证(SASL/SCRAM 和 SASL/PLAIN)。
订阅模式
订阅模式指定了流如何选择要使用的 Kafka 主题。 支持三种模式:
| 模式 | Description | 示例: |
|---|---|---|
subscribe |
以逗号分隔的主题名称列表 | KafkaSubscriptionMode(subscribe="topic1,topic2") |
subscribe_pattern |
Java正则表达式模式匹配主题名称 | KafkaSubscriptionMode(subscribe_pattern="events-.*") |
assign |
用于指定主题-分区分配的 JSON | KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}') |
架构配置
使用 JSON 架构 格式定义消息键和值的结构。 对于 Kafka 源,payload_schema 对应于 Kafka 消息值(即 Kafka 键值模型中的 value),而 key_schema 对应于 Kafka 消息键。 必须至少提供 payload_schema 或 key_schema 其中之一。
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string"}'
' }'
'}'
)
),
key_schema=SchemaConfig(
json_schema='{"type": "string"}'
),
)
如果没有为密钥或有效负载提供架构,则它被视为一个简单的字符串。
引入和回填
该 ingestion_config 参数配置如何捕获和存储流数据以供训练和服务。
对流的访问受引入表控制:
- 对引入表的
SELECT权限将授予对流的读取访问权限。 - 对引入表的
MANAGE权限将授予删除访问权限。
有关表权限的更多信息,请参阅 表 和 Unity Catalog 权限参考。
摄取管道
创建流时,Databricks 将启动托管引入管道,该管道会持续读取 Kafka 主题中的消息,并将其写入 Delta 表(引入表)。 管道从最新的 Kafka 偏移量开始,并持续运行,仅捕获在创建流后到达的新消息。 此引入表用于使用流式处理特征进行训练。 删除流后,也会删除其引入管道和引入表。
摄取目标
ingestion_destination 指定用于写入流数据的三部分 Delta 表名称。
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)
数据引入表架构
引入表包含消息数据以及元数据列:
| 列 | 类型 | Description |
|---|---|---|
key |
各不相同(从 key_schema 起) |
Kafka 消息密钥,根据提供的架构进行结构化。 |
value |
各不相同(从 payload_schema 起) |
Kafka 消息值(有效负载),根据提供的架构进行结构化。 |
stream_record_timestamp |
TIMESTAMP |
记录的时间戳。 对于前向填充数据,这是 Kafka 代理的引入时间戳。 回填数据由客户提供。 |
kafka_topic |
STRING |
作为所用记录来源的 Kafka 主题。 |
kafka_partition |
INT |
作为所用记录来源的 Kafka 分区。 |
kafka_offset |
LONG |
记录在其所属分区中的 Kafka 偏移量。 |
record_source |
STRING |
"stream"(从实时 Kafka 流中前向填充)或 "backfill"(来自回填源)。 |
回填源
由于前向填充管道是从最新的 Kafka 偏移量开始,因此无法捕获该流创建之前已存在的消息。 若要为训练提供历史数据覆盖范围,请配置可选的回填源。
配置回填源后,Databricks 会运行一次性 MERGE INTO 作业,该作业使用 record_source="backfill" 将回填行复制到引入表中。 MERGE 仅在重叠检查器确认回填源和前填充流具有重叠时间戳之后才会运行(请参阅 回填和实时流数据之间的重叠)。 如果在 2 天内未满足重叠条件,则 MERGE 仍会运行以避免无限期阻塞。
回填表必须包含一个名为 stream_record_timestamp、类型为 TIMESTAMP 且采用 UTC 时区的列。 其他 Kafka 元数据列(kafka_topic、kafka_partition、kafka_offset)如果在回填源中存在,则会原样传递;否则将其设为 NULL。
from databricks.feature_engineering.entities import StreamBackfillSource
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
backfill_source=StreamBackfillSource(
delta_table_name="my_catalog.my_schema.historical_events"
),
)
回填和实时流数据之间的重叠
在回填表和引入表之间运行 MERGE 之前,会先进行重叠检查,比较这两个表中的时间戳:
-
回填最大值:回填源中的最大值
stream_record_timestamp。 -
引入最小值:引入表中的最小
stream_record_timestamp行数(record_source="stream")。
只有回填的最新时间戳至少比引入表中最早的时间戳晚 1 小时,MERGE 才会继续执行。 这种重叠机制可确保引入表中没有间隙。 如果在 2 天内未满足重叠条件,则 MERGE 仍会运行以避免无限期阻塞。
由于引入管道从最新的 Kafka 偏移量开始读取,因此只能捕获在创建流之后到达的消息。 回填源必须包含延伸到引入时间范围内的数据,而不只是到流创建时间为止。
例如,如果你在下午 3:00 创建流,则前向填充管道会从下午 3:00 起开始读取消息。 回填源必须包含时间戳截至至少下午4:00的数据(即比前向填充开始时间晚1小时),才能通过重叠检查。 这意味着你应在下午4:00之后更新回填表,以确保摄取表不会出现空档。
重复数据删除
使用 deduplication_columns 指定列路径,以便在引入过程中识别回填数据与向前填充流数据之间的重复行。 对嵌套字段使用点表示法(例如 "value.user_id")。
根据您的数据选择去重列:
- 如果流中的每个记录都包含唯一标识符(例如),
value.transaction_id请使用该列进行重复数据删除。 - 如果回填源包含
kafka_partition和kafka_offset列,请使用这些列来唯一标识每条记录。 - 如果未指定去重列,则默认去重键为
key、value和stream_record_timestamp的完整组合。 不建议这样做,因为这种严格的条件匹配很容易导致重复项。
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)
管理流
获取流
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
列出流
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)
将 include_schemas=True 设置为包含完整的架构模式详情。 架构可能很庞大,因而导致操作长时间运行。 若要单独检索架构,请使用 get_stream。
删除流
删除流也会删除其引入管道和引入表。
Warning
引用已删除流的任何模型或功能将不再有权访问基础流数据。 如果需要此数据,但不再需要流,请在删除之前创建引入表的副本。
client.delete_stream(name="my_catalog.my_schema.my_stream")
示例笔记本
有关创建流、定义流功能并部署到服务终结点的端到端示例,请参阅以下笔记本: