数据分享:数据发布与数据订阅

从源工作区发布源数据库中的表,并在目标工作区订阅该发布。发布和订阅发生在不同工作区,分别返回各自的资源和状态。

任务流程

  1. 在源工作区读取可发布的表,并由调用方选择发布范围。

  2. 在源工作区创建发布,得到发布资源和发布信息。

  3. 在目标工作区读取订阅列表,找到对应的订阅资源。

  4. 使用订阅资源订阅发布,并根据返回状态和目标数据库确认结果。

准备

需要的内容

在本页中的作用

源工作区和目标工作区

发布和订阅各自的资源作用域。

源数据库

读取可发布的表并建立发布。

表范围

选择发布全部表或指定表。

目标工作区

接收发布邀请并创建订阅。

发布名称、订阅名称和备注

标识发布和目标工作区中的订阅。

发布范围有两种明确选择:发布全部表,或发布一个或多个指定表。选择指定表时,先从源数据库的可发布表列表中选择;选择全部表时,不再提供单独的表列表。

创建发布只表示源工作区已收到创建请求,不表示目标工作区已经订阅。

发布并订阅数据

下面的示例先确认指定表位于可发布列表中,再创建发布并在目标工作区找到订阅资源。最后订阅该发布,并检查订阅结果。

import moi_product_sdk as sdk

def publish_and_subscribe(
    source_workspace,
    target_workspace,
    source_database_id,
    source_table_id,
    target_workspace_id,
    publish_name,
    subscription_name,
    remark,
):
    source_tables = source_workspace.data_share_source_tables(source_database_id)
    if source_table_id not in {item.id for item in source_tables.list}:
        raise ValueError("source table is not available for publication")
    table_scope = sdk.DataShareScope(
        mode="selected",
        object_ids=[source_table_id],
    )
    publish, created = source_workspace.create_data_share_publish(
        publish_name, source_database_id, table_scope,
        [target_workspace_id], sdk.with_data_share_remark(remark),
    )
    if publish.id != created.id:
        raise RuntimeError("publication handle and result do not match")

    subscriptions = target_workspace.data_share().subscriptions()
    subscription_id = next(
        (item.id for item in subscriptions.list if item.pub_name == created.name),
        "",
    )
    if not subscription_id:
        raise RuntimeError("subscription invitation was not found")
    subscription = target_workspace.data_share_subscription(subscription_id)
    subscribed = subscription.subscribe(subscription_name)
    if subscribed.status != "subscribed":
        raise RuntimeError(f"unexpected subscription status: {subscribed.status}")
    return subscribed
import (
	"context"
	"fmt"

	sdk "github.com/matrixorigin/matrixflow/sdk/go-sdk"
)

func publishAndSubscribe(
	ctx context.Context,
	sourceWorkspace, targetWorkspace *sdk.WorkspaceHandle,
	sourceDatabaseID, sourceTableID, targetWorkspaceID, publishName, subscriptionName, remark string,
) error {
	sourceTables, err := sourceWorkspace.DataShareSourceTables(ctx, sourceDatabaseID)
	if err != nil {
		return err
	}
	found := false
	for _, table := range sourceTables.GetList() {
		if table.GetId() == sourceTableID {
			found = true
			break
		}
	}
	if !found {
		return fmt.Errorf("source table %q is not available for publication", sourceTableID)
	}

	tableScope := &sdk.DataShareScope{
		Mode:      "selected",
		ObjectIds: []string{sourceTableID},
	}
	publish, created, err := sourceWorkspace.CreateDataSharePublish(
		ctx, publishName, sourceDatabaseID, tableScope,
		[]string{targetWorkspaceID}, sdk.WithDataShareRemark(remark),
	)
	if err != nil {
		return err
	}
	if publish.ID() != created.GetId() {
		return fmt.Errorf("publication handle and result do not match")
	}

	subscriptions, err := targetWorkspace.DataShare().Subscriptions(ctx)
	if err != nil {
		return err
	}
	subscriptionID := ""
	for _, item := range subscriptions.GetList() {
		if item.GetPubName() == created.GetName() {
			subscriptionID = item.GetId()
			break
		}
	}
	if subscriptionID == "" {
		return fmt.Errorf("subscription invitation for %q was not found", created.GetName())
	}
	subscription, err := targetWorkspace.DataShareSubscription(subscriptionID)
	if err != nil {
		return err
	}
	subscribed, err := subscription.Subscribe(ctx, subscriptionName)
	if err != nil {
		return err
	}
	if subscribed.GetStatus() != "subscribed" {
		return fmt.Errorf("unexpected subscription status: %s", subscribed.GetStatus())
	}
	return nil
}

示例中的指定表来自源数据库的可发布表列表。创建发布后,返回的发布资源与发布信息指向同一次创建结果;可以使用发布资源更新范围、目标工作区或备注,也可以删除发布。

结果确认

订阅结果包含当前订阅状态和目标数据库。需要确认后续状态变化时,在目标工作区重新读取订阅列表;不要将提交订阅或取消订阅的请求成功当作最终状态。

限制

数据分享会影响跨工作区的数据可见性。创建、更新、删除发布,或订阅、取消订阅前,确认源数据库、表范围、目标工作区和访问权限。跨工作区邀请出现和订阅状态变更需要通过目标工作区的订阅列表确认。

下一步

最后更新于