Lakeflow 管道中的流会将数据写入流式表或物化视图。 以下示例演示如何定义默认流、将流与其目标分开定义、将多个 Kafka 主题的数据写入流式表、运行一次性回填,以及用追加流处理替换 UNION 查询。
有关流概述,请参阅 使用 Lakeflow 管道流以增量方式加载和处理数据。
示例:创建默认流
创建管道时,通常会定义表或视图以及支持它的查询。 例如,此查询通过从 customers_silver 读取来创建一个名为 customers_bronze 的流式表。 流式表及其默认流程将在单一步骤中同时创建。
SQL
CREATE OR REFRESH STREAMING TABLE customers_silver
AS SELECT * FROM STREAM(customers_bronze)
Python
from pyspark import pipelines as dp
@dp.table()
def customers_silver():
return spark.readStream.table("customers_bronze")
流表的默认流是 append 流,每次更新时都会添加新行,其名称与目标名称相同。 这是使用管道(在单个步骤中创建流及其目标)的最常用方法,你可以使用它引入或转换数据。 有关流概念的详细信息,请参阅 使用 Lakeflow 管道流以增量方式加载和处理数据。
示例:定义流与其目标分开
还可以为单独定义的表创建流。 结果与创建默认流相同,包括对流表和流使用相同的名称:
Python
from pyspark import pipelines as dp
# create streaming table
dp.create_streaming_table("customers_silver")
# add a flow
@dp.append_flow(
target = "customers_silver")
def customer_silver():
return spark.readStream.table("customers_bronze")
SQL
-- create a streaming table
CREATE OR REFRESH STREAMING TABLE customers_silver;
-- add a flow
CREATE FLOW customers_silver
AS INSERT INTO customers_silver BY NAME
SELECT * FROM STREAM(customers_bronze);
通过独立于其目标定义流,可以创建多个流,这些流将数据追加到同一目标。 在 Python 接口中使用 @dp.append_flow 修饰器,或在 SQL 接口中使用 CREATE FLOW...INSERT INTO 子句,以便为如下任务添加流:
- 添加将数据追加到现有流式处理表的流源,而无需完全刷新。 例如,你可能有一个表,该表合并了你在其中运行的每个区域的区域数据。 新区域推出后,无需执行完全刷新即可将新区域数据添加到表中。 请参阅 示例:将多个 Kafka 主题的数据写入流式表。
- 通过追加缺失的历史数据(回填)来更新流表。 您可以使用
INSERT INTO ONCE语法创建仅运行一次的历史数据回填。 请参阅 示例:使用管道运行一次性数据回填 和 回填历史数据。 - 将来自多个源的数据整合并写入一个流处理表,以替代在查询中使用
UNION子句。 使用追加流处理,而不是UNION允许以增量方式更新目标表,而无需运行 完全刷新更新。 请参阅 示例:使用追加流处理而不是UNION。
对于 Python 查询,请使用 create_streaming_table() 函数创建目标表。
Important
- 如果需要使用 预期定义数据质量约束,请将目标表的预期定义为函数的
create_streaming_table()一部分或现有表定义。 不能在@append_flow定义中定义期望。 - 流由 流名称标识,此名称用于标识流式处理检查点。 使用流名称标识检查点意味着:
- 如果管道中的现有流已重命名,则检查点不会进行传递,并且重命名的流实际上是一个全新的流。
- 不能重用管道中的流名称,因为现有检查点与新流定义不匹配。
示例:从多个 Kafka 主题写入流式处理表
以下示例创建一个名为 kafka_target 的流式处理表,并从两个 Kafka 主题写入该流式处理表:
Python
from pyspark import pipelines as dp
dp.create_streaming_table("kafka_target")
# Kafka stream from multiple topics
@dp.append_flow(target = "kafka_target")
def topic1():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,...")
.option("subscribe", "topic1")
.load()
)
@dp.append_flow(target = "kafka_target")
def topic2():
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,...")
.option("subscribe", "topic2")
.load()
)
SQL
CREATE OR REFRESH STREAMING TABLE kafka_target;
CREATE FLOW
topic1
AS INSERT INTO
kafka_target BY NAME
SELECT * FROM
read_kafka(bootstrapServers => 'host1:port1,...', subscribe => 'topic1');
CREATE FLOW
topic2
AS INSERT INTO
kafka_target BY NAME
SELECT * FROM
read_kafka(bootstrapServers => 'host1:port1,...', subscribe => 'topic2');
若要详细了解 read_kafka() SQL 查询中使用的表值函数,请参阅 SQL 语言参考中的 read_kafka 。
在 Python 中,可以编程方式创建面向单个表的多个流。 以下示例显示了 Kafka 主题列表的此模式。
注释
此模式的要求与使用 for 循环创建表的要求相同。 必须将 Python 值显式传递给定义流的函数。 请参阅 在for 循环中创建表。
from pyspark import pipelines as dp
dp.create_streaming_table("kafka_target")
topic_list = ["topic1", "topic2", "topic3"]
for topic_name in topic_list:
@dp.append_flow(target = "kafka_target", name=f"{topic_name}_flow")
def topic_flow(topic=topic_name):
return (
spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,...")
.option("subscribe", topic)
.load()
)
示例:运行一次性数据回填
如果要运行查询以将数据追加到现有流式处理表,请使用 append_flow。
追加一组现有数据后,有多个选项:
- 如果希望查询在新数据抵达回填目录时能够自动添加,请保持查询有效。
- 如果希望这是一次性回填,并且永远不会再次运行,请在运行管道一次后删除查询。
- 如果希望查询仅运行一次,并且仅在数据被完全刷新时再次运行,请在追加流程中将
once参数设置为True。 在 SQL 中,使用INSERT INTO ONCE。
以下示例运行查询以将历史数据追加到流式处理表:
Python
from pyspark import pipelines as dp
@dp.table()
def csv_target():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format","csv")
.load("path/to/sourceDir")
@dp.append_flow(
target = "csv_target",
once = True)
def backfill():
return spark.read
.format("cloudFiles")
.option("cloudFiles.format","csv")
.load("path/to/backfill/data/dir")
SQL
CREATE OR REFRESH STREAMING TABLE csv_target
AS SELECT * FROM
read_files(
"path/to/sourceDir",
"csv"
);
CREATE FLOW
backfill
AS INSERT INTO ONCE
csv_target BY NAME
SELECT * FROM
read_files(
"path/to/backfill/data/dir",
"csv"
);
有关更深入的示例,请参阅 使用管道回填历史数据。
示例:使用追加流处理而不是 UNION
可以使用追加流查询来合并多个源并写入单个流式表,而不是使用带有UNION子句的查询。 使用追加流查询替代UNION,可以在不进行完全刷新的情况下,从多个源追加到流式处理表中。
以下 Python 示例包含一个查询,该查询将多个数据源与子句组合在一起 UNION :
@dp.create_table(name="raw_orders")
def unioned_raw_orders():
raw_orders_us = (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/us")
)
raw_orders_eu = (
spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/eu")
)
return raw_orders_us.union(raw_orders_eu)
以下示例将 UNION 查询替换为追加流查询:
Python
dp.create_streaming_table("raw_orders")
@dp.append_flow(target="raw_orders")
def raw_orders_us():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/us")
@dp.append_flow(target="raw_orders")
def raw_orders_eu():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/eu")
# Additional flows can be added without the full refresh that a UNION query would require:
@dp.append_flow(target="raw_orders")
def raw_orders_apac():
return spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "csv")
.load("/path/to/orders/apac")
SQL
CREATE OR REFRESH STREAMING TABLE raw_orders;
CREATE FLOW
raw_orders_us
AS INSERT INTO
raw_orders BY NAME
SELECT * FROM
STREAM read_files(
"/path/to/orders/us",
format => "csv"
);
CREATE FLOW
raw_orders_eu
AS INSERT INTO
raw_orders BY NAME
SELECT * FROM
STREAM read_files(
"/path/to/orders/eu",
format => "csv"
);
-- Additional flows can be added without the full refresh that a UNION query would require:
CREATE FLOW
raw_orders_apac
AS INSERT INTO
raw_orders BY NAME
SELECT * FROM
STREAM read_files(
"/path/to/orders/apac",
format => "csv"
);
示例:使用 transformWithState 监测传感器心跳
以下示例演示一个从 Kafka 读取并验证传感器是否定期发出心跳信号的有状态处理器。 如果在 5 分钟内未收到心跳信号,处理器会向目标 Delta 表提交一条记录以供分析。
有关生成自定义有状态应用程序的详细信息,请参阅 生成自定义有状态应用程序。
注释
RocksDB 是从 Databricks Runtime 17.2 开始的默认状态提供程序。 如果查询因不支持的提供程序异常而失败,请添加必要的管道配置,执行系统的完全刷新或检查点重置,然后重新运行管道:
"configuration": {
"spark.sql.streaming.stateStore.providerClass": "com.databricks.sql.streaming.state.RocksDBStateStoreProvider",
"spark.sql.streaming.stateStore.rocksdb.changelogCheckpointing.enabled": "true"
}
from typing import Iterator
import pandas as pd
from pyspark import pipelines as dp
from pyspark.sql.functions import col, from_json
from pyspark.sql.streaming import StatefulProcessor, StatefulProcessorHandle
from pyspark.sql.types import StructType, StructField, LongType, StringType, TimestampType
KAFKA_TOPIC = "<your-kafka-topic>"
output_schema = StructType([
StructField("sensor_id", LongType(), False),
StructField("sensor_type", StringType(), False),
StructField("last_heartbeat_time", TimestampType(), False)])
class SensorHeartbeatProcessor(StatefulProcessor):
def init(self, handle: StatefulProcessorHandle) -> None:
# Define state schema to store sensor information (sensor_id is the grouping key)
state_schema = StructType([
StructField("sensor_type", StringType(), False),
StructField("last_heartbeat_time", TimestampType(), False)])
self.sensor_state = handle.getValueState("sensorState", state_schema)
# State variable to track the previously registered timer
timer_schema = StructType([StructField("timer_ts", LongType(), False)])
self.timer_state = handle.getValueState("timerState", timer_schema)
self.handle = handle
def handleInputRows(self, key, rows, timerValues) -> Iterator[pd.DataFrame]:
# Process one row from input and update state
pdf = next(rows)
row = pdf.iloc[0]
# Store or update the sensor information in state using current timestamp
current_time = pd.Timestamp(timerValues.getCurrentProcessingTimeInMs(), unit='ms')
self.sensor_state.update((
row["sensor_type"],
current_time
))
# Delete old timer if already registered
if self.timer_state.exists():
old_timer = self.timer_state.get()[0]
self.handle.deleteTimer(old_timer)
# Register a timer for 5 minutes from current processing time
expiry_time = timerValues.getCurrentProcessingTimeInMs() + (5 * 60 * 1000)
self.handle.registerTimer(expiry_time)
# Store the new timer timestamp in state
self.timer_state.update((expiry_time,))
# No output on input processing, output only on timer expiry
return iter([])
def handleExpiredTimer(self, key, timerValues, expiredTimerInfo) -> Iterator[pd.DataFrame]:
# Emit output row based on state store
if self.sensor_state.exists():
state = self.sensor_state.get()
output = pd.DataFrame({
"sensor_id": [key[0]], # Use grouping key as sensor_id
"sensor_type": [state[0]],
"last_heartbeat_time": [state[1]]
})
# Remove the entry for the sensor from the state store
self.sensor_state.clear()
# Remove the timer state entry
self.timer_state.clear()
yield output
def close(self) -> None:
pass
dp.create_streaming_table("sensorAlerts")
# Define the schema for the Kafka message value
sensor_schema = StructType([
StructField("sensor_id", LongType(), False),
StructField("sensor_type", StringType(), False),
StructField("sensor_value", LongType(), False)])
@dp.append_flow(target = "sensorAlerts")
def kafka_delta_flow():
return (
spark.readStream
.format("kafka")
.option("subscribe", KAFKA_TOPIC)
.option("startingOffsets", "earliest")
.load()
.select(from_json(col("value").cast("string"), sensor_schema).alias("data"), col("timestamp"))
.select("data.*", "timestamp")
.withWatermark('timestamp', '1 hour')
.groupBy(col("sensor_id"))
.transformWithStateInPandas(
statefulProcessor = SensorHeartbeatProcessor(),
outputStructType = output_schema,
outputMode = 'update',
timeMode = 'ProcessingTime'))