使用同步表提供湖仓数据

通过同步表,您可以使用 Lakebase Postgres 为 Lakehouse 数据提供服务。 Unity 目录表同步到 Postgres,以便应用程序可以直接以低延迟查询 Lakehouse 数据。 此过程通常称为反向 ETL。 Lakehouse 针对分析和扩充进行优化,而 Lakebase 专为需要快速查找样式查询和事务一致性的操作工作负荷而设计。

显示从 Lakehouse 到 Lakebase 到应用程序的数据流的体系结构图

什么是同步表?

借助同步表,您可以通过 Lakebase Postgres 提供来自 Unity Catalog 的分析级别数据,使这些数据可供需要低延迟查询和完整的 ACID 事务支持的应用程序使用。 它们通过确保数据可在实时应用程序中使用,弥合分析存储与操作系统之间的差距。

支持的源

已同步的表格支持以下 Unity Catalog 源类型:

  • 托管表和外部 Delta 表
  • 托管表和外部 Iceberg 表
  • 视图和具体化视图

工作原理

Databricks 同步表会在 Lakebase 中创建 Unity Catalog 数据的托管副本。 创建同步表时,将得到:

  1. Unity Catalog 中引用同步管道的同步表
  2. Lakebase 中的 Postgres 表 (应用程序可查询的只读表)

显示同步表中三表关系的关系图

例如,可以将黄金表、工程特征或 ML 输出从 analytics.gold.user_profiles 同步到新的同步表 analytics.gold.user_profiles_synced 中。 在 Postgres 中,Unity 目录架构名称将成为 Postgres 架构名称,因此如下所示 gold.user_profiles_synced

SELECT * FROM gold.user_profiles_synced WHERE user_id = 12345;

应用程序使用标准 Postgres 驱动程序进行连接,并查询同步的数据以及它们自身的运行状态。

警告

尽管可以直接在 Postgres 中修改同步表,但Azure Databricks严格建议只运行读取查询来保护源的数据完整性。 有关同步表支持的操作,请参阅 Postgres 中已同步表上允许的操作

同步管道使用托管的 Lakeflow 管道,根据源表中的更改持续更新 Unity Catalog 同步表和 Postgres 表。 每个同步最多可以使用 16 个到 Lakebase 数据库的连接。

Lakebase Postgres 支持多达 1,000 个具有事务保证的并发连接,因此应用程序可以读取扩充的数据,同时处理同一数据库中的插入、更新和删除。

同步模式

根据应用程序需求选择正确的同步模式:

模式 说明 何时使用 性能
快照 所有数据的一次性副本 源数据变更每个周期 > 10% 的行,或源数据不支持 CDF(视图、Iceberg 表) 如果修改 >10% 源数据,则效率更高 10 倍
引发 按需或间隔运行的计划更新 源数据行按已知频率发生变化。 每次刷新都会传播插入、更新和删除操作。 良好的成本/滞后平衡。 如果以 <5 分钟为间隔运行,成本较高
连续的 秒级延迟的实时流传输 更改必须几乎实时地显示在 Lakebase 中 滞后时间最低,成本最高。 最小 15 秒间隔

触发模式和连续模式要求在源表上启用变更数据源 (CDF)。 如果未启用 CDF,则会在 UI 中看到一条警告,其中包含要运行的确切 ALTER TABLE 命令。 有关更改数据馈送的更多详细信息,请参阅 在 Databricks 上使用 Delta Lake 更改数据馈送

注释

不支持 CDF(如视图、具体化视图和 Iceberg 表)的源只能在 快照 模式下同步。 对于快照模式,源必须支持 SELECT *

示例用例

可以将同步表用于数据服务用例,例如:

  • 向 Databricks 应用传送新用户档案的个性化引擎
  • 提供湖屋中计算的模型预测或特征值的应用程序
  • 面向客户的仪表板,实时展示 KPI
  • 为立即行动提供风险评分的欺诈检测服务
  • 支持从湖屋数据中提供增强型客户记录的工具

创建同步表

先决条件

您需要:

  • 启用了 Lakebase 的 Databricks 工作区。
  • Lakebase 项目(请参阅 创建项目)。
  • 一个待同步的 Unity Catalog 表。
  • 创建同步表的权限。 你需要在使用的任何模式上具有USE_SCHEMACREATE_TABLE权限。

对于触发式连续式模式,源表上必须启用变更数据馈送 (CDF):

ALTER TABLE your_catalog.your_schema.your_table
SET TBLPROPERTIES (delta.enableChangeDataFeed = true)

有关容量规划和数据类型兼容性,请参阅 数据类型和兼容性容量规划

UI

  1. 转到工作区边栏中的 目录 ,然后选择要同步的 Unity 目录表。

    Catalog Explorer 中显示所选表

  2. 从表详细信息视图中单击“ 创建>同步表 ”。

    “创建”按钮的下拉菜单中显示“已同步表”选项

  3. “创建同步表 ”对话框中:

    目录和架构列表仅包括当前用户具有 USE_SCHEMACREATE_TABLE 权限的 Unity 目录架构。 如果未看到所需的架构,请与目录管理员确认权限。

    1. 表名称:输入已同步表的名称(它与源表在同一目录和架构中创建)。 这会同时创建 Unity 目录同步表和可以查询的 Postgres 表。

    2. 数据库类型:选择 Lakebase 无服务器(自动缩放)。

    3. 同步模式:根据需求选择 快照触发模式或 连续 模式(请参阅上面的 同步模式 )。

    4. 配置项目、分支和数据库选择。

    5. 验证 主键 是否正确(通常自动检测到)。

      重要

      同步表中主键的列不允许为空。 主键列中为 null 的行 将从同步中排除

    6. (可选)如果源表中的两行可能共享相同的主键,请选择时间序列键来配置去重。 指定时间序键时,同步表仅包含每个主键具有最新时间序键值的行。 有关没有时间序列键的故障模式,请参阅 重复键

    如果选择了“触发”或“连续”模式,但尚未启用“更改数据馈送”,则会看到一条警告,其中包含要运行的确切命令。 有关数据类型兼容性问题,请参阅 数据类型和兼容性

    单击“ 创建 ”以创建同步表。

  4. 监控目录中的同步表。 “ 概述 ”选项卡显示同步状态、配置、管道状态和上次同步时间戳。 立即同步以进行手动刷新。

CLI

databricks postgres create-synced-table my-catalog.sales.orders \
  --json '{
    "spec": {
      "source_table_full_name": "main.sales.orders",
      "branch": "projects/my-project/branches/production",
      "primary_key_columns": ["order_id"],
      "scheduling_policy": "SNAPSHOT",
      "postgres_database": "mydb",
      "create_database_objects_if_missing": true
    }
  }'

SYNCED_TABLE_ID位置参数使用格式 catalog.schema.table。 在 Postgres 中,表{table}是在模式{schema}中创建的,位于你用postgres_database设置的数据库中(此处为mydb)。 该命令会等待操作默认完成。 有关所有可用选项,请参阅 databricks postgres create-synced-table

Python SDK

from databricks.sdk import WorkspaceClient
from databricks.sdk.service.postgres import (
    SyncedTable,
    SyncedTableSyncedTableSpec,
    SyncedTableSyncedTableSpecSyncedTableSchedulingPolicy,
)

w = WorkspaceClient()

synced_table = w.postgres.create_synced_table(
    synced_table=SyncedTable(spec=SyncedTableSyncedTableSpec(
        source_table_full_name="main.sales.orders",
        branch="projects/my-project/branches/production",
        primary_key_columns=["order_id"],
        scheduling_policy=SyncedTableSyncedTableSpecSyncedTableSchedulingPolicy.SNAPSHOT,
        postgres_database="mydb",
        create_database_objects_if_missing=True,
    )),
    synced_table_id="my-catalog.sales.orders",
).wait()

print(f"Synced table created: {synced_table.name}")

synced_table_id 使用 catalog.schema.table 格式,并成为 Unity Catalog 的同步表名称。 在 Postgres 中,表{table}是在模式{schema}中创建的,位于你用postgres_database设置的数据库中(此处为mydb)。

Java SDK

import com.databricks.sdk.WorkspaceClient;
import com.databricks.sdk.service.postgres.*;
import java.util.List;

WorkspaceClient w = new WorkspaceClient();

SyncedTable syncedTable = w.postgres().createSyncedTable(
    new CreateSyncedTableRequest()
        .setSyncedTableId("my-catalog.sales.orders")
        .setSyncedTable(new SyncedTable()
            .setSpec(new SyncedTableSyncedTableSpec()
                .setSourceTableFullName("main.sales.orders")
                .setBranch("projects/my-project/branches/production")
                .setPrimaryKeyColumns(List.of("order_id"))
                .setSchedulingPolicy(SyncedTableSyncedTableSpecSyncedTableSchedulingPolicy.SNAPSHOT)
                .setPostgresDatabase("mydb")
                .setCreateDatabaseObjectsIfMissing(true))))
    .waitForCompletion();

System.out.println("Synced table created: " + syncedTable.getName());

curl

curl -X POST "https://your-workspace.cloud.databricks.com/api/2.0/postgres/synced_tables?synced_table_id=my-catalog.sales.orders" \
  -H "Authorization: Bearer ${DATABRICKS_TOKEN}" \
  -H "Content-Type: application/json" \
  -d '{
    "spec": {
      "source_table_full_name": "main.sales.orders",
      "branch": "projects/my-project/branches/production",
      "primary_key_columns": ["order_id"],
      "scheduling_policy": "SNAPSHOT",
      "postgres_database": "mydb",
      "create_database_objects_if_missing": true
    }
  }'

这将返回一个长时间运行的操作。 轮询返回的 name 字段,直到 done: true。 请参阅 长时间运行的操作。 有关身份验证设置,请参阅 “身份验证”。

计划或触发后续同步

初始快照在创建时自动运行。 对于 快照触发 模式,必须明确触发后续同步。 连续 模式是自我管理。

数据库表同步流程任务

Lakeflow 作业中的数据库表同步管道任务将同步表的管道作为工作流步骤运行。 使用表更新触发器或计划来配置任务。

源表更新时触发

当源 Unity Catalog 表更新时触发该作业。 在触发式模式下,仅增量应用新变更,可在避免“连续”模式始终在线成本的同时,提供近乎实时的数据新鲜度。

  1. 在边栏中,单击“ 工作流”。
  2. 单击“ 创建作业 ”或打开现有作业。
  3. 在“ 任务 ”选项卡上,单击“ + 添加其他任务类型”。
  4. 引入和转换下,选择 “数据库表同步”管道
  5. “管道 ”字段中,选择与同步表关联的管道。
  6. “计划和触发器”下,单击“ 添加触发器”。
  7. 选择 “表更新 ”作为触发器类型。
  8. “表”下,选择要监视的源 Unity 目录表。
  9. 单击“ 保存”。

按计划触发

按固定频率执行同步。 非常适合 快照 模式,其中夜间或每周完全刷新通常是最有效的模式。

  1. 按照上述步骤 1-5 将 数据库表同步管道 任务添加到作业。
  2. “计划和触发器”下,单击“ 添加触发器”。
  3. 选择 已计划 作为触发类型。
  4. 设置 cron 计划和时区,然后单击“ 保存”。

检查同步状态

检查同步表的当前状态和上次同步时间:

UI

“目录”中,导航到同步表并选择“ 概述 ”选项卡。它显示当前同步状态、管道状态和上次同步时间戳。

Python SDK

from databricks.sdk import WorkspaceClient

w = WorkspaceClient()

table = w.postgres.get_synced_table("synced_tables/my-catalog.sales.orders")
print(f"State: {table.status.detailed_state}")
print(f"Last sync: {table.status.last_sync_time}")
print(f"Message: {table.status.message}")

Java SDK

import com.databricks.sdk.WorkspaceClient;
import com.databricks.sdk.service.postgres.SyncedTable;

WorkspaceClient w = new WorkspaceClient();

SyncedTable table = w.postgres().getSyncedTable("synced_tables/my-catalog.sales.orders");
System.out.println("State: " + table.getStatus().getDetailedState());
System.out.println("Last sync: " + table.getStatus().getLastSyncTime());
System.out.println("Message: " + table.getStatus().getMessage());

curl

curl "https://your-workspace.cloud.databricks.com/api/2.0/postgres/synced_tables/my-catalog.sales.orders" \
  -H "Authorization: Bearer ${DATABRICKS_TOKEN}"

数据类型和兼容性

创建同步表时,Unity 目录数据类型将映射到 Postgres 类型。 复杂类型(ARRAY、MAP、结构)以 JSONB 形式存储在 Postgres 中。

源列类型 Postgres 列类型
BIGINT BIGINT
BINARY BYTEA
BOOLEAN BOOLEAN
DATE DATE
DECIMAL(p, s) NUMERIC
DOUBLE 双精度
FLOAT real
INT INTEGER
INTERVAL INTERVAL
SMALLINT SMALLINT
STRING TEXT
TIMESTAMP 时间戳与时区
TIMESTAMP_NTZ 无时区时间戳
TINYINT SMALLINT
ARRAY<元素类型> JSONB
MAP<keyType,valueType> JSONB
STRUCT<fieldName:fieldType[, ...]> JSONB

注释

不支持 GEOGRAPHY、GEOMETRY、VARIANT 和 OBJECT 类型。

处理无效字符

在 Unity Catalog 的 STRING、ARRAY、MAP 或 STRUCT 列中,某些字符(如 null 字节(0x00))是允许的,但在 Postgres 的 TEXT 或 JSONB 列中则不支持。 这可能会导致同步失败,并出现如下错误:

ERROR: invalid byte sequence for encoding "UTF8": 0x00
ERROR: unsupported Unicode escape sequence DETAIL: \u0000 cannot be converted to text
  • 第一个错误发生在顶级字符串列中出现 null 字节时,该列直接映射到 Postgres TEXT
  • 当嵌套在复杂类型(STRUCTARRAYMAP)内部的字符串中出现空字节时,会发生第二个错误,该字符串会被序列化为 JSONB。 在序列化期间,所有字符串都会被转换为 Postgres TEXT 格式,其中 \u0000 是不允许的。

解决方案:

  • 清理字符串字段:在同步之前删除不受支持的字符。 对于“STRING”列中的null字节:

    SELECT REPLACE(column_name, CAST(CHAR(0) AS STRING), '') AS cleaned_column FROM your_table
    
  • 转换为 BINARY:对于需要保留原始字节的 STRING 列,请转换为 BINARY 类型。

容量规划

规划同步表实现时,请考虑以下资源要求:

  • 连接使用情况:每个同步表最多使用 16 个到 Lakebase 数据库的连接,这些连接数计入实例的连接限制。
  • 大小配额:所有同步表的总逻辑数据配额为 16 TB。 如果需要更大的配额,请联系 Databricks 支持部门。 单个表没有配额,但对于需要刷新的表,Databricks 建议不超过 1 TB。
  • 完全刷新大小:触发完全刷新时,在新同步完成之前不会删除 Postgres 中的旧版本。 这两个版本在刷新期间暂时计入逻辑数据库大小配额。
  • 每个源的表数量:单个源表最多可同步到 20 个表。
  • 命名要求:数据库、架构和表名只能包含字母数字字符和下划线([A-Za-z0-9_]+)。
  • 源标识符指南:避免在源 Unity 目录表中的列名或表名中使用大写字母或特殊字符。 如果保留这些标识符,则必须在 Postgres 中引用这些标识符。
  • 架构演变:触发模式和连续模式仅支持累加架构更改(如添加列)。
  • 重复键:如果两行在源表中具有相同的主键,则同步管道将失败,除非使用 时间顺序键配置重复数据删除。
  • API 幂等性:同步表 API 接口具有幂等性,因此在出现瞬时错误时请重试,以确保操作及时完成。
  • 更新速率:对于 Lakebase 自动扩展,同步管道支持“连续”和“触发式”写入,速度约为每容量单位 (CU) 每秒 150 行;支持“快照”写入,速度高达每容量单位 (CU) 每秒 2,000 行。

Postgres 中同步表上允许的操作

Azure Databricks建议仅在 Postgres 中为同步表执行以下操作,以防止意外覆盖或数据不一致:

  • 只读查询
  • 创建索引
  • 删除表(从 Unity Catalog 中删除此同步表以释放空间)

虽然可以通过其他方式修改 Postgres 中的同步表,但它会干扰同步管道。

所有权和权限

同步表由内部 databricks_writer_<dbid> 角色(而不是创建该表的用户)拥有,因为同步管道管理它(请参阅 Postgres 角色)。 仅所有者命令(例如配置行级别安全性)不能直接在同步表上运行。

注释

这是对 Postgres 一般规则的一个例外:如果你的 Azure Databricks 身份的登录名在 Postgres 中作为角色存在,则你自己创建的对象归你的 Azure Databricks 身份所有。 该管道会为你创建同步表。

创建同步表的用户的访问权限

创建同步表时,系统会自动向您的 Azure Databricks 身份授予使用该表的权限。 无需 databricks_superuser 执行任何操作。 您的身份已被授予以下针对同步表的权限:

Object 特权 Purpose
同步表 SELECTDELETETRUNCATE 读取或清除表格
Schema USAGECREATE 使用架构并创建索引等对象

你未获得 INSERTUPDATE。 数据管道拥有该表的数据,因此直接写入的数据会在下次刷新时被覆盖。 DELETETRUNCATE 只会清除表格。 下一次刷新会从源重新填充表。

此访问权限源自您在该同步表上的 Unity Catalog 权限,并由 Unity Catalog 管理。 若要更改它,请更新用户的 Unity 目录权限。 你不能在 Postgres 中直接通过 Azure Databricks 标识REVOKE它。

注释

此访问权限与创建该同步表的身份相关联。 更改管道的 运行身份 不会将其重新分配。 若要使用不同的所有者标识,请以该标识重新创建同步表。

管理同步表访问

创建同步表后, databricks_superuser 可以从 Postgres 读取同步表。 databricks_superuser 具有 pg_read_all_data,这使得该角色可以读取所有表。 它还具有 pg_write_all_data 权限,允许此角色写入所有表。 这意味着 databricks_superuser 也可以写入 Postgres 中的同步表。 如果你需要在目标表中进行紧急更改,Lakebase 支持这种写入行为。 但是,Azure Databricks建议改为在源表中进行修复。

  • databricks_superuser 还可以向其他用户授予这些权限。

    GRANT USAGE ON SCHEMA synced_table_schema TO user;
    
    GRANT SELECT ON synced_table_name TO user;
    
  • databricks_superuser 可以撤销这些特权:

    REVOKE USAGE ON SCHEMA synced_table_schema FROM user;
    
    REVOKE {SELECT | INSERT | UPDATE | DELETE} ON synced_table_name FROM user;
    

管理同步表操作

databricks_superuser 可以管理哪些用户有权对同步表执行特定操作。 同步表支持的操作包括:

  • CREATE INDEX
  • ALTER INDEX
  • DROP INDEX
  • DROP TABLE

已同步表的所有其他 DDL 操作都被拒绝。

若要向其他用户授予这些权限, databricks_superuser 必须先创建 databricks_auth以下扩展:

CREATE EXTENSION IF NOT EXISTS databricks_auth;

然后, databricks_superuser 可以添加用户来管理同步表:

SELECT databricks_synced_table_add_manager('"synced_table_schema"."synced_table"'::regclass, '[user]');

databricks_superuser 可以将用户从同步表的管理中删除。

SELECT databricks_synced_table_remove_manager('[table]', '[user]');

databricks_superuser 可以查看所有管理员:

SELECT * FROM databricks_synced_table_managers;

删除同步表

从 Unity 目录中删除同步表也会删除相应的 Postgres 表。

UI

“目录”中,找到已同步的表,单击 “Kebab”菜单图标。 菜单,然后选择“ 删除”。

Python SDK

from databricks.sdk import WorkspaceClient

w = WorkspaceClient()

w.postgres.delete_synced_table("synced_tables/my-catalog.sales.orders").wait()

Java SDK

import com.databricks.sdk.WorkspaceClient;

WorkspaceClient w = new WorkspaceClient();

w.postgres().deleteSyncedTable("synced_tables/my-catalog.sales.orders").waitForCompletion();

curl

curl -X DELETE "https://your-workspace.cloud.databricks.com/api/2.0/postgres/synced_tables/my-catalog.sales.orders" \
  -H "Authorization: Bearer ${DATABRICKS_TOKEN}"

了解详细信息

任务 说明
创建项目 创建 Lakebase 项目
连接到数据库 了解 Lakebase 的连接选项
在 Unity 目录中注册数据库 使 Lakebase 数据在 Unity 目录中可见,以便统一治理和跨源查询
Unity 目录集成 了解治理和权限

目录集成

  • 目录重复: 在目标为 Postgres 数据库的标准目录中创建同步表,该数据库也注册为单独的数据库目录会导致同步表显示在标准目录和数据库目录下的 Unity 目录中。

其他选项

若要将数据同步到非 Databricks 系统,请参阅 Partner Connect 反向 ETL 解决方案 ,例如人口普查或 Hightouch。