Streaming Delta Live Tables流表引用STREAM()语法疑问
Delta Live Tables中
STREAM()函数的正确使用规则 认为所有流式DLT表(Streaming Live Table)引用时必须包裹STREAM()函数的认知存在偏差。STREAM()的核心作用是声明对目标表采用增量流式读取语义,和被引用表本身是流式表还是批式表没有强制绑定关系。
官方示例的逻辑说明
你看到的无报错写法是DLT的标准合法语法,对应建表逻辑如下:
CREATE OR REFRESH STREAMING LIVE TABLE sales_orders_cleaned( CONSTRAINT valid_order_number EXPECT (order_number IS NOT NULL) ON VIOLATION DROP ROW ) COMMENT "The cleaned sales orders with valid_order_number(s) and partitioned by order_datetime." AS SELECT f.customer_id, f.customer_name, f.number_of_line_items, timestamp(from_unixtime((cast(f.order_datetime as long)))) as order_datetime, date(from_unixtime((cast(f.order_datetime as long)))) as order_date, f.order_number, f.ordered_products, c.state, c.city, c.lon, c.lat, c.units_purchased, c.loyalty_segment FROM STREAM(LIVE.sales_orders_raw) f LEFT JOIN LIVE.customers c ON c.customer_id = f.customer_id AND c.customer_name = f.customer_name
STREAM(LIVE.sales_orders_raw):将sales_orders_raw标记为整个流查询的驱动流源,管道每次触发计算时,仅处理该表自上次成功运行后新写入的增量数据,不会重复扫描全量历史。- 未包裹
STREAM()的LIVE.customers:即便该表本身是通过cloud_files定义的流式DLT表,这里的读取语义为每次触发计算时,拉取customers表当前时刻的最新全量快照,作为维度表和新流入的订单数据做关联。计算过程不会追踪customers表的历史变更,也不会因为customers表写入新数据触发计算,属于DLT原生支持的流-静态关联场景。
示例涉及的两个上游输入表均为流式DLT表,不影响上述读取逻辑,表定义如下:
CREATE OR REFRESH STREAMING LIVE TABLE sales_orders_raw COMMENT "The raw sales orders, ingested from /databricks-datasets." AS SELECT * FROM cloud_files("/databricks-datasets/retail-org/sales_orders/", "json", map("cloudFiles.inferColumnTypes", "true")) CREATE OR REFRESH STREAMING LIVE TABLE customers COMMENT "The customers buying finished products, ingested from /databricks-datasets." AS SELECT * FROM cloud_files("/databricks-datasets/retail-org/customers/", "csv");
语法规则总结
STREAM()的使用只和查询需要的读取语义有关,和被引用表本身的定义类型无关:- 需要对某张表做增量流式读取,即把它作为流源、每次只消费新写入的增量数据时,才需要用
STREAM()包裹表名。一个合法的DLT流查询至少需要包含一个被STREAM()修饰的流源。 - 只需要读取某张表当前时刻的最新快照,把它作为静态维度表、参考表参与计算时,无论该表是批式DLT表还是流式DLT表,都不需要加
STREAM()。
- 需要对某张表做增量流式读取,即把它作为流源、每次只消费新写入的增量数据时,才需要用
- 只有做流-流关联时(即关联两侧的表都作为流源,任意一侧有新数据写入都要触发匹配计算),才需要给两侧的表都加上
STREAM()包裹。
内容的提问来源于stack exchange,提问作者Minura Punchihewa
相关产品推荐
相关产品推荐

