使用独立流式表

独立流式表是指在 Lakeflow 管道之外定义、注册到 Unity Catalog 中,并额外支持流式或增量数据处理的表。 系统会自动为每个流式处理表创建一个管道。 可以使用流式处理表从 Kafka 和云对象存储进行增量数据加载。

您可以从 Databricks SQL 仓库或在无服务器通用计算上运行的笔记本中创建和刷新独立的流式表。 有关两个计算选项之间的差异的详细信息,请参阅 独立管道的要求

若要使用笔记本中的Python创建和刷新独立流式处理表,请参阅将Python与独立管道配合使用

注释

若要了解如何将 Delta Lake 表用作流式处理源和接收器,请参阅 Delta Lake 表流式处理读取和写入

要求

有关创建、刷新和查询独立流式处理表的计算选项、权限和其他要求,请参阅 独立管道的要求

创建流式处理表

流式处理表由 Databricks SQL 中的 SQL 查询定义。 创建流式处理表时,源表中当前的数据用于生成流式处理表。 之后,通常按照预定的时间表刷新表,以从源表中获取任何新增数据,并将其附加到流式表中。

当你创建流式处理表时,你将被视为该表的所有者。

若要从现有表创建流式处理表,请使用该语句 CREATE STREAMING TABLE,如以下示例所示:

CREATE OR REFRESH STREAMING TABLE sales
  SCHEDULE EVERY 1 hour
  AS SELECT product, price FROM STREAM raw_data;

在这种情况下,流式表 sales 是从 raw_data 表的特定列创建的,并计划每小时刷新一次。 使用的查询必须是 流式处理 查询。 要使用流式处理语义从源中读取,请使用 STREAM 关键字。

用于刷新的计算

使用 CREATE OR REFRESH STREAMING TABLE 语句创建流表时,初始数据刷新和填充会立即开始。 这些操作不使用 Databricks SQL 仓库计算。 相反,流式表依赖无服务器管道进行创建和刷新。 系统会自动为每个流式处理表创建和管理专用无服务器管道。

使用自动加载程序加载文件

若要从卷中的文件创建流数据表,请使用 Auto Loader。 对云对象存储中的大多数数据引入任务使用自动加载程序。 Auto Loader 和管道旨在以增量且幂等的方式,在数据到达云存储时实时加载不断增长的数据。

若要在 Databricks SQL 中使用自动加载程序,请使用函数 read_files 。 以下示例演示如何使用自动加载程序将 JSON 文件卷读取到流式处理表中:

CREATE OR REFRESH STREAMING TABLE sales
  SCHEDULE EVERY 1 hour
  AS SELECT * FROM STREAM read_files(
    "/Volumes/my_catalog/my_schema/my_volume/path/to/data",
    format => "json"
  );

若要从云存储中读取数据,还可以使用自动加载程序:

CREATE OR REFRESH STREAMING TABLE sales
  SCHEDULE EVERY 1 hour
  AS SELECT *
  FROM STREAM read_files(
    'abfss://myContainer@myStorageAccount.dfs.core.windows.net/analysis/*/*/*.json',
    format => "json"
  );

若要了解自动加载程序,请参阅什么是自动加载程序? 若要了解有关在 SQL 中使用自动加载程序的详细信息,请参阅 从对象存储加载数据

从其他源流式处理引入

有关从其他源(包括 Kafka)引入的示例,请参阅 在管道中加载数据

使用 Auto CDC 流应用变更数据捕获 (CDC)

使用 FLOW AUTO CDC 子句将源中的变更数据捕获 (CDC) 记录处理到流式表中。 以前,该 MERGE INTO 语句通常用于处理 Azure Databricks 上的 CDC 记录。 但是, MERGE INTO 由于序列外记录或需要复杂的逻辑重新排序记录,可能会生成不正确的结果。 请参阅 变更数据捕获和快照

AUTO CDC 通过自动处理无序记录来简化 CDC。 指定用于标识记录的键、排序的序列列,以及是将结果存储为 SCD 类型 1(直接更新)还是 SCD 类型 2(历史记录跟踪)。

以下示例创建了一个流式表,该表使用 SCD 类型 1 应用 CDC 变更:

CREATE OR REFRESH STREAMING TABLE target
  FLOW AUTO CDC
  FROM stream(cdc_data.users)
  KEYS (userId)
  SEQUENCE BY sequenceNum
  STORED AS SCD TYPE 1;

以下示例使用 SCD 类型 2 来保留更改历史记录:

CREATE OR REFRESH STREAMING TABLE target
  FLOW AUTO CDC
  FROM stream(cdc_data.users)
  KEYS (userId)
  APPLY AS DELETE WHEN operation = "DELETE"
  SEQUENCE BY sequenceNum
  COLUMNS * EXCEPT (operation, sequenceNum)
  STORED AS SCD TYPE 2;

有关自动 CDC 选项和行为的完整详细信息,请参阅 AUTO CDC API:使用管道简化更改数据捕获。 有关完整的语法参考,请参阅 CREATE STREAMING TABLE

使用 REPLACE WHERE 流程进行选择性批量替换

使用 FLOW REPLACE WHERE 子句可重新计算并覆盖流式表中的目标子集,而无需重新处理整个表的历史记录。 REPLACE WHERE 流非常适合用于联接和聚合、后期到达数据、上游重新处理、架构演变和回填的增量批处理。

有关 REPLACE WHERE 流的完整详细信息(包括要求、谓词覆盖和增量刷新),请参阅独立流式表的 REPLACE WHERE 流

使用 REPLACE USING 流应用部分快照替换

重要

REPLACE USING 流仍处于 Beta 阶段

使用 FLOW REPLACE USING 子句使流式表与部分快照流保持同步。 每次更新时,REPLACE USING 流程会替换所有匹配指定键列的行,其他行保持不变。 SEQUENCE BY 列会对更新进行排序,因此即使更新到达顺序混乱,也始终应用键下序列号最大的更新。 例如:

CREATE OR REFRESH STREAMING TABLE payments_current
FLOW REPLACE USING (payment_id) SEQUENCE BY payment_date BY NAME
SELECT payment_id, booking_id, status, payment_date
FROM STREAM(samples.wanderbricks.payments);

BY NAME 必需。 它根据列名而非位置进行匹配。

REPLACE USING 在独立流式表中的行为与在 Lakeflow 管道中的行为相同。 关于其工作原理、执行顺序、预期、限制和示例,请参见 使用 REPLACE USING 流程进行部分快照替换。 以下差异适用于独立流式表:

  • 用SQL定义流程。CREATE OR REFRESH STREAMING TABLE 中使用内联 SQL FLOW REPLACE USING 子句编写 REPLACE USING 流。 独立的 CREATE FLOW 语句是 Lakeflow 管道构造,不适用于独立的流式表。
  • 计算资源由系统代为管理。 独立流式表在系统托管的无服务器管道上运行,且要求使用 Databricks Runtime 18.2 及更高版本。 你不能在经典和无服务器计算之间做选择。

只引入新数据

默认情况下,该 read_files 函数在创建表期间读取源文件夹中的所有现有数据,然后在每次刷新时处理新到达的记录。

若要避免在创建表时源文件夹中已存在的引入数据,请将 includeExistingFiles 选项设置为 false。 这意味着只有在创建表之后到达文件夹中的数据才会被处理。 例如:

CREATE OR REFRESH STREAMING TABLE sales
  SCHEDULE EVERY 1 hour
  AS SELECT *
  FROM STREAM read_files(
    '/path/to/files',
    includeExistingFiles => false
  );

运行时版本

流表始终运行在最新的Databricks SQL运行时版本上。 pipelines.channel 表属性此前用于选择 previewcurrent 运行时通道,现已不再受支持,且不会生效。 如果现有定义包含了该属性,则可以安全忽略,无需移除。

隐藏敏感数据

可以使用流式处理表来对访问表的用户隐藏敏感数据。 一种方法是定义查询,以便完全排除敏感列或行。 或者,可以根据查询用户的权限应用列掩码或行筛选器。 例如,可以为不在 tax_id 组中的用户隐藏 HumanResourcesDept 列。 为此,在创建流式处理表期间使用 ROW FILTERMASK 语法。 有关详细信息,请参阅 行筛选器和列掩码

刷新流式处理表

流式处理表会自动创建和使用无服务器管道来处理刷新操作。 刷新由管道管理,更新由用于创建流数据表的 Databricks SQL 仓库进行监控。 流式表可通过在计划上运行的管道进行更新。

即使已计划刷新,也可以随时调用手动刷新。 刷新由与流式处理表一起自动创建的同一管道处理。

刷新流式处理表:

REFRESH STREAMING TABLE sales;

可以使用DESCRIBE TABLE EXTENDED检查最新刷新状态。

注释

在使用按时间顺序查看查询之前,可能需要刷新流式处理表。

若要了解如何计划刷新,请参阅 计划刷新。 计划刷新可发送更新通知,并且您可以为刷新设置性能模式

刷新的工作原理

流式表刷新仅评估自上次更新后到达的新行,并仅追加新数据。

每次刷新都使用流式处理表的当前定义来处理此新数据。 修改流式处理表定义不会自动重新计算现有数据。 如果修改与现有数据不兼容(例如更改数据类型),则下一次刷新将失败并显示错误。

以下示例说明流式处理表定义的更改如何影响刷新行为:

  • 删除筛选器不会重新处理以前筛选的行。
  • 更改列投影不会影响处理现有数据的方式。
  • 具有静态快照的联接使用初始处理时的快照状态。 那些本应与更新后的快照匹配但延迟到达的数据将被忽略。 如果维度延迟到达,可能会导致事实丢弃。
  • 修改现有列的 CAST 会导致错误。

如果现有的流式处理表无法支持数据更改的方式,你可以执行完整刷新。

完全刷新流式处理表

全量刷新会根据最新定义重新处理源中所有可用的数据。 对于不保留数据全部历史记录或保留期较短的源(如 Kafka),不建议调用完全刷新,因为完全刷新会截断现有的数据。 如果数据在源中不再可用,则可能无法恢复旧数据。

例如:

REFRESH STREAMING TABLE sales FULL;

计划和监视刷新

您可以按计划自动刷新流式表,也可以在上游数据发生变化时自动刷新,还可以配置刷新超时、通知和性能模式。 请参阅 刷新计划

控制对流式处理表的访问

流式表支持多样化的访问控制,以支持数据共享,同时避免泄露潜在的敏感数据。 具有 MANAGE 权限的流式处理表所有者或用户可以向其他用户授予 SELECT 权限。 具有 SELECT 流式传输表访问权限的用户无需 SELECT 访问该流式传输表所引用的表。 此访问控制支持数据共享,同时控制对基础数据的访问。

您还可以修改流式传输表的所有者。

授予对流式处理表的权限

若要授予对流式处理表的访问权限,请使用 GRANT 语句

GRANT <privilege_type> ON <st_name> TO <principal>;

privilege_type可以是:

  • SELECT - 用户可以选择 (SELECT) 流式处理表。
  • REFRESH - 用户可以选择 (REFRESH) 流式处理表。 刷新是使用所有者的权限运行的。

以下示例创建一个流式表单,并向用户授予 SELECT 和刷新权限:

CREATE OR REFRESH STREAMING TABLE st_name AS SELECT * FROM source_table;

-- Grant read-only access:
GRANT SELECT ON st_name TO read_only_user;

-- Grant read and refresh access:
GRANT SELECT ON st_name TO refresh_user;
GRANT REFRESH ON st_name TO refresh_user;

有关授予 Unity 目录安全对象特权的详细信息,请参阅 Unity 目录特权参考

撤销对流式处理表的权限

若要撤销对流式处理表的访问权限,请使用 REVOKE 语句

REVOKE privilege_type ON <st_name> FROM principal;

当撤销了流式处理表所有者或已被授予流式处理表 SELECTMANAGE 权限的任何其他用户对源表的 SELECT 权限,或者删除了源表时,流式处理表所有者或被授予访问权限的用户仍然可以查询该流式处理表。 但是,会发生以下行为:

  • 流式处理表所有者或失去流式处理表访问权限的其他人无法再 REFRESH 使用该流式处理表,流式处理表会随着时间推移而过时。
  • 如果计划自动执行,则下一个计划 REFRESH 失败或未运行。

以下示例从 SELECT 撤销了 read_only_user 权限:

REVOKE SELECT ON st_name FROM read_only_user;

更改流式处理表的所有者

对独立流式表具有 MANAGE 权限的用户可以通过目录资源管理器设置新的所有者。 新所有者可以是用户本人,也可以是用户拥有服务主体用户角色的服务主体。

  1. 在 Azure Databricks 工作区中,单击 “数据”图标。目录以打开目录资源管理器。

  2. 选择要更新的流式处理表。

  3. 在右侧边栏的关于此流式传输表下,找到所有者,然后单击 铅笔图标。 编辑。

    注释

    如果你收到一则消息,提示你通过在管道设置中更改 Run as 用户来更新所有者,则该流式表是在 Lakeflow 管道中定义的,而不是定义为独立表。 该消息包含一个指向管道设置的链接,您可以在其中更改运行用户。

  4. 为流式传输表选择新的所有者。

    所有者对其拥有的流式传输表自动拥有 MANAGESELECT 权限。 如果您将服务主体设置为您所拥有的流式表的所有者,而您在该流式表上未明确拥有 SELECTMANAGE 权限,则此更改将导致您完全失去对该流式表的访问权限。 在这种情况下,系统会提示显式地提供这些权限。

    同时选中授予 MANAGE授予 SELECT 权限,然后单击保存

  5. 单击“ 保存 ”以更改所有者。

流式处理表的所有者已更新。 所有将来的更新都使用新所有者的标识运行。

当所有者丧失源表权限时

如果您更改了所有者,而新所有者无权访问源表(或底层源表上的 SELECT 权限已被撤销),用户仍可查询该流式表。 但是:

  • 他们无法 REFRESH 流式表。
  • 流式表的下一次计划刷新将失败。

失去对源数据的访问权限会阻止更新,但不会立即阻止对现有流式表的读取。

从流式处理表永久删除记录

重要

对流式处理表的 REORG 语句的支持现推出公共预览版

注释

  • 对流式处理表使用 REORG 语句需要 Databricks Runtime 15.4 及更高版本。
  • 尽管可以将 REORG 语句与任何流式处理表一起使用,但只有在从启用了 删除矢量 的流式处理表中删除记录时,才需要使用该语句。 在未启用删除向量的情况下与流式处理表一起使用时,该命令不起作用。

若要从启用了删除矢量的流式处理表的基础存储中物理删除记录(例如为了 GDPR 合规),必须执行其他步骤,以确保 VACUUM 操作在流式处理表的数据上运行。

若要从基础存储中物理删除记录,请执行以下作:

  1. 更新记录或删除流数据表中的记录。
  2. 在流式处理表上运行REORG语句,并指定APPLY (PURGE)参数。 例如,REORG TABLE <streaming-table-name> APPLY (PURGE);
  3. 等待流式处理表的数据保留期传递。 默认数据保留期为 7 天,但可以使用表属性进行配置 delta.deletedFileRetentionDuration 。 请参阅为“按时间顺序查看”查询配置数据保留
  4. REFRESH 流式处理表。 请参阅刷新流式处理表。 在 REFRESH 操作的 24 小时内,包括 VACUUM 操作在内的管道维护任务(其中 VACUUM 用于确保记录被永久删除)会自动运行。

使用查询历史记录监视运行

可以使用查询历史记录页访问查询详细信息和查询分析,这些信息有助于识别在用于运行流式处理表更新的管道中表现不佳的查询和瓶颈。 有关查询历史记录和查询配置文件中可用的信息的概述,请参阅 查询历史记录查询配置文件

重要

此功能目前以公共预览版提供。 工作区管理员可以从 预览 页控制对此功能的访问。 请参阅 管理 Azure Databricks 预览版

与流式处理表相关的所有语句都显示在查询历史记录中。 可以使用“语句”下拉列表筛选器来选择任何命令并检查相关查询。 所有 CREATE 语句后都跟随一个在管道上异步执行的 REFRESH 语句。 这些 REFRESH 语句通常包括详细的查询计划,用于提供优化性能的见解。

若要访问 REFRESH 查询历史记录 UI 中的语句,请使用以下步骤:

  1. 单击 “历史记录”图标。 在左侧栏中打开 “查询历史记录 ”UI。
  2. 从“语句”下拉列表筛选器中选择 REFRESH 复选框。
  3. 单击查询语句的名称可查看摘要详细信息,例如查询持续时间和聚合指标。
  4. 单击“ 查看查询配置文件 ”打开查询配置文件。 有关浏览查询配置文件的详细信息,请参阅查询配置文件
  5. (可选)可以使用“查询源”部分中的链接打开相关的查询或管道。

还可以使用 SQL 编辑器中的链接或附加到 SQL 仓库的笔记本访问查询详细信息。

从外部客户端访问流式处理表

若要从不支持开放 API 的外部 Delta Lake 或 Iceberg 客户端访问流式处理表,可以使用 兼容性模式。 兼容性模式会创建流式表的只读版本,任何 Delta Lake 或 Iceberg 客户端均可访问该版本。

其他资源