连接器、导入与导出任务¶
在工作区中创建连接器,并用它提交文件导入或导出任务。创建调用会返回任务资源;用该资源读取任务详情、运行记录或状态,再决定是否调整、重试或删除。
任务流程¶
调用方在目标工作区选择连接器配置和传输任务的输入。
SDK 创建或绑定连接器,并用它提交导入任务或导出任务。
SDK 返回对应的任务资源。
使用这个资源读取任务详情、运行记录、任务文件或运行状态。
准备¶
需要的内容 |
在本页中的作用 |
|---|---|
已认证的客户端和已选择的工作区 |
确定连接器和任务的资源作用域。 |
连接器名称、来源类型和使用类型 |
创建连接器。 |
导入或导出配置 |
定义数据来源、目标位置和处理方式。 |
目标卷或待导出文件 |
指定传输的落点或内容。 |
外部系统的访问条件 |
让连接器能够访问调用方选定的外部系统。 |
导入和导出配置由调用方提供。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
结果确认¶
导入任务可读取详情、运行记录和任务文件;导出任务可读取详情、任务文件和运行状态。重试、暂停、恢复、更新或删除前,先检查当前任务信息和数据影响。
限制¶
创建任务表示请求已提交,不表示数据已完成导入或导出。
导出可能向外部存储写入数据。创建、重试或删除前,确认目标位置、访问条件和覆盖策略。
重试前检查失败原因和数据是否已部分写入,避免重复处理。