连接器、导入与导出任务

在工作区中创建连接器,并用它提交文件导入或导出任务。创建调用会返回任务资源;用该资源读取任务详情、运行记录或状态,再决定是否调整、重试或删除。

任务流程

  1. 调用方在目标工作区选择连接器配置和传输任务的输入。

  2. SDK 创建或绑定连接器,并用它提交导入任务或导出任务。

  3. SDK 返回对应的任务资源。

  4. 使用这个资源读取任务详情、运行记录、任务文件或运行状态。

准备

需要的内容

在本页中的作用

已认证的客户端和已选择的工作区

确定连接器和任务的资源作用域。

连接器名称、来源类型和使用类型

创建连接器。

导入或导出配置

定义数据来源、目标位置和处理方式。

目标卷或待导出文件

指定传输的落点或内容。

外部系统的访问条件

让连接器能够访问调用方选定的外部系统。

导入和导出配置由调用方提供。SDK 不会从文件路径、显示名称或本地环境推断来源、目标或覆盖策略。

创建文件导入任务

下面的函数创建连接器并提交文件导入任务。提交成功只表示任务创建请求已返回,不表示数据已经导入。

import moi_product_sdk as sdk


def create_file_import_task(
    workspace,
    connector_name,
    source_type,
    usage_type,
    config_type,
    task_name,
    volume_id,
    source_uris,
    load_mode_config,
    file_filter_config,
):
    connector, _ = workspace.create_connector(
        connector_name, source_type, usage_type
    )
    task, _ = connector.create_file_import_task(
        sdk.FileImportTaskSpec(
            config_type,
            source_uris,
            load_mode_config,
            file_filter_config,
            name=task_name,
            volume_id=volume_id,
        )
    )
    return connector, task
import (
	"context"

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

func createFileImportTask(
	ctx context.Context,
	workspace *sdk.WorkspaceHandle,
	connectorName string,
	sourceType, usageType, configType int,
	taskName, volumeID string,
	sourceURIs []string,
	loadModeConfig, fileFilterConfig map[string]any,
) (*sdk.ConnectorHandle, *sdk.ImportTaskHandle, error) {
	connector, _, err := workspace.CreateConnector(ctx, connectorName, sourceType, usageType)
	if err != nil {
		return nil, nil, err
	}
	task, _, err := connector.CreateFileImportTask(ctx, sdk.FileImportTaskSpec{
		ConfigType:       configType,
		Name:             taskName,
		VolumeID:         volumeID,
		URIs:             sourceURIs,
		LoadModeConfig:   loadModeConfig,
		FileFilterConfig: fileFilterConfig,
	})
	if err != nil {
		return nil, nil, err
	}
	return connector, task, nil
}

保留函数返回的导入任务资源。遇到网络超时后,先查询已知任务,而不是直接重复提交。

检查导入任务

使用上一步返回的导入任务资源读取详情、运行记录和任务文件。它们分别反映本次读取到的任务信息、运行记录和文件列表;不要将创建请求返回写成传输完成。

details = import_task.info()
runs = import_task.runs()
files = import_task.files()
details, err := importTask.Info(ctx)
if err != nil {
	return err
}
runs, err := importTask.Runs(ctx)
if err != nil {
	return err
}
files, err := importTask.Files(ctx)
if err != nil {
	return err
}
_ = details
_ = runs
_ = files

创建导出任务

导出任务需要任务名称、导出类型、至少一种目标配置以及至少一个待导出的文件。文件的资源标识和完整路径必须来自已知资源或调用方选择。可以使用上一步返回的连接器资源创建导出任务。

export_task, _ = connector.create_export_task_from_spec(
    sdk.ExportTaskSpec(
        "export-orders",
        "oss",
        2,
        sdk.ExportTaskConfig(s3_config={"path": "/exports/orders"}),
        [sdk.ExportTaskFile("file-1", ["sales", "orders"], is_raw=True)],
    )
)
isRaw := true
exportTask, _, err := connector.CreateExportTaskFromSpec(ctx, sdk.ExportTaskSpec{
	Name:          "export-orders",
	ConnectorName: "oss",
	Type:          2,
	Config:         sdk.ExportTaskConfig{S3Config: map[string]any{"path": "/exports/orders"}},
	Files: []sdk.ExportTaskFile{{
		FileID:   "file-1",
		FullPath: []string{"sales", "orders"},
		IsRaw:    &isRaw,
	}},
})
if err != nil {
	return err
}
_ = exportTask

结果确认

导入任务可读取详情、运行记录和任务文件;导出任务可读取详情、任务文件和运行状态。重试、暂停、恢复、更新或删除前,先检查当前任务信息和数据影响。

限制

  • 创建任务表示请求已提交,不表示数据已完成导入或导出。

  • 导出可能向外部存储写入数据。创建、重试或删除前,确认目标位置、访问条件和覆盖策略。

  • 重试前检查失败原因和数据是否已部分写入,避免重复处理。

下一步

最后更新于