Copula Lab 与现有数据仓库同步专家标注数据的集成方案

文章导读
专家标注数据要从 Copula Lab 进入现有数据仓库,关键不是找现成连接器,而是先把“导出字段”和“仓库表字段”的对齐关系定死,再用定时任务做增量拉取,最后用数据库本身的冲突处理能力保证重复执行不产生脏数据。Copula Lab 当前没有公开的原生数据仓库连接器,通用的做法是走它提供的导出接口或人工导出文件,由外部 ETL 脚本统一拉取。下面给出可直接照用的同步任务配置模板、字段映射表、冲突处
📋 目录
  1. 梳理导出字段与数仓表结构的映射
  2. 编写定时拉取脚本并写入暂存区
  3. 使用 SQL 完成增量合并与幂等控制
  4. 配置同步日志和失败告警
  5. 执行一次全量比对验证数据一致性
A A

专家标注数据要从 Copula Lab 进入现有数据仓库,关键不是找现成连接器,而是先把“导出字段”和“仓库表字段”的对齐关系定死,再用定时任务做增量拉取,最后用数据库本身的冲突处理能力保证重复执行不产生脏数据。Copula Lab 当前没有公开的原生数据仓库连接器,通用的做法是走它提供的导出接口或人工导出文件,由外部 ETL 脚本统一拉取。下面给出可直接照用的同步任务配置模板、字段映射表、冲突处理规则,以及验证比对 SQL。

适用场景是 Copula Lab 支持导出标注结果文件或提供 REST API,且数据仓库允许通过 SQL 脚本写入。操作动作是固定导出字段顺序、用外部脚本分页拉取、写入暂存区、再用 MERGE 或 ON CONFLICT 做幂等合并。验证方式是执行全量 COUNT 和关键字段 SUM 比对。风险边界是字段类型和 ID 生成规则必须提前确认,否则增量合并会因主键不一致产生重复记录。

梳理导出字段与数仓表结构的映射

先导出一份最小样本,把 Copula Lab 侧的字段名、导出顺序、类型和数仓目标表字段一一列出来。不要依赖字段名自动匹配,因为两边对同一概念命名可能不同,比如 Copula Lab 里叫 document_id,数仓里叫 doc_id。建议先建一张固定的映射表,后续所有脚本都从这张表读取对应关系,避免在代码里散落多处硬编码字段。

映射表至少要包含:源字段名、导出列序号、目标字段名、目标表字段类型、转换规则。常见类型转换规则是:字符串统一转为 VARCHAR 并去掉首尾空白;日期时间先确认是 UTC 还是本地时区,再转换为目标库的 TIMESTAMP;标注结果是 JSON 结构的,建议存成原生 JSON 类型或拆成子表,不要直接塞进长文本字段。

-- 数仓侧建表示例,字段按映射表定义
CREATE TABLE dw.label_data (
    id           BIGINT PRIMARY KEY,
    doc_id       VARCHAR(64)  NOT NULL,
    labeler_id   VARCHAR(64),
    label_type   VARCHAR(32),
    label_value  JSON,
    annotated_at TIMESTAMP,
    updated_at   TIMESTAMP,
    sync_batch   VARCHAR(64)
);

映射确认后,把源字段和目标字段的对应关系写成一份 CSV 或数据库配置表,脚本启动时读取。这样后续 Copula Lab 导出格式变动时,只需更新配置表,不用改主要代码。

编写定时拉取脚本并写入暂存区

同步通道分两步:拉取文件或调用接口,然后写入暂存区。不要在拉取过程中直接写目标表,避免中间失败污染正式数据。暂存区可以先落在本地 CSV,也可以直接写入数仓的 staging 表。下面给出一个 Python 脚本骨架,使用分页参数拉取接口,写入 CSV:

Copula Lab 与现有数据仓库同步专家标注数据的集成方案
import requests, csv, time

API_URL = "https://your-copula-instance/api/exports"
TOKEN = "your_token"
PAGE_SIZE = 500

session = requests.Session()
session.headers.update({"Authorization": f"Bearer {TOKEN}"})

with open("staging_label_data.csv", "w", newline="") as f:
    writer = csv.writer(f)
    writer.writerow(["id", "doc_id", "labeler_id", "label_type", "label_value", "annotated_at"])
    page = 1
    while True:
        resp = session.get(API_URL, params={"page": page, "size": PAGE_SIZE})
        resp.raise_for_status()
        data = resp.json()
        if not data:
            break
        for item in data:
            writer.writerow([
                item["id"],
                item["doc_id"],
                item.get("labeler_id"),
                item.get("label_type"),
                item.get("label_value"),
                item.get("annotated_at")
            ])
        page += 1
        time.sleep(0.5)  # 避免过快请求
print("拉取完成")

如果数据量不大,也可以用 Shell 脚本配合 curl 拉取,再用 jq 解析 JSON 写入文件。脚本要记录本次拉取批次号,例如用日期时间生成 sync_batch,方便追溯。定时调度建议用系统自带的 cron 或任务平台,不要在脚本内部写死循环。

使用 SQL 完成增量合并与幂等控制

拉取完成后,把暂存区数据导入数仓的暂存表,再通过幂等合并写入最终表。幂等的核心是确定主键,通常用 Copula Lab 导出的记录 ID 作为主键。如果同一记录被重复拉取,后来的数据要覆盖之前的数据,而不是插入两条。使用 MERGEINSERT ... ON CONFLICT 可以满足这个需求。

-- PostgreSQL 示例,使用 ON CONFLICT
INSERT INTO dw.label_data (
    id, doc_id, labeler_id, label_type, label_value, annotated_at, updated_at, sync_batch
)
SELECT t.id, t.doc_id, t.labeler_id, t.label_type,
       t.label_value::jsonb, t.annotated_at, NOW(), t.sync_batch
FROM staging.label_data AS t
ON CONFLICT (id) DO UPDATE SET
    labeler_id  = EXCLUDED.labeler_id,
    label_type  = EXCLUDED.label_type,
    label_value = EXCLUDED.label_value,
    annotated_at = EXCLUDED.annotated_at,
    updated_at  = NOW(),
    sync_batch  = EXCLUDED.sync_batch;

如果是 SQL Server 或 MySQL,MySQL 用 INSERT ... ON DUPLICATE KEY UPDATE,SQL Server 用 MERGE。关键是消除重复插入:先判断主键是否已存在,存在则更新,不存在才插入。为了性能,暂存表导入前先按主键去重,避免同一批次内重复记录引起冲突。

配置同步日志和失败告警

同步任务失败不能只靠人工发现。建议在每次同步任务的关键步骤写日志,日志至少记录字段要包含:任务时间、批次号、拉取记录数、写入暂存区行数、合并成功行数、失败错误信息、耗时。这些日志可以写入数据库日志表,也可以输出到文件供监控工具采集。

Copula Lab 与现有数据仓库同步专家标注数据的集成方案
-- 日志表结构示例
CREATE TABLE dw.sync_log (
    id BIGSERIAL PRIMARY KEY,
    task_name VARCHAR(64),
    batch_id VARCHAR(64),
    started_at TIMESTAMP,
    finished_at TIMESTAMP,
    pulled_count INT,
    merged_count INT,
    error_message TEXT
);

告警触发命令要放在脚本末尾。例如如果返回码非零或合并前后行数差异超过阈值,就发送通知到工作群。通常做法是脚本捕获异常后调用 webhook 地址:

# 失败时发送告警
if [ $? -ne 0 ]; then
  curl -X POST -H "Content-Type: application/json" \
    -d '{"text": "Copula Lab 同步任务失败,批次号: $BATCH_ID"}' \
    https://your-notification-webhook.example
fi

告警阈值要结合环境确认,不要用拍脑袋的固定值。重点是让失败在下一个同步周期开始前被发现。

执行一次全量比对验证数据一致性

配置好同步后,做一次全量比对,确认迁入数仓的数据没有丢数。比对方法是分别对源端和目标端执行 COUNT 和 SUM,比较两个值是否一致。源端的统计通常来自 Copula Lab 的导出接口返回的元数据,如果接口没有提供总数,就只能对导出的源文件做统计。

-- 目标端统计
SELECT COUNT(*) AS total_rows,
       SUM(id) AS id_sum
FROM dw.label_data;

-- 源端统计(假设源文件已落到本地表)
SELECT COUNT(*) AS total_rows,
       SUM(id) AS id_sum
FROM staging.source_snapshot;

如果 COUNT 一致但 SUM 不一致,说明有字段读取错位;如果 COUNT 也不一致,说明有数据没有拉到或合并时被去重过度。还要抽查几条关键字段,比如取 doc_idlabel_value,在源和目标中分别按主键取对应值比较。全量比对不是只做一次,首次部署后和每次源导出格式变动后都应该执行。比对脚本可以复用,放入定时任务中,但不要影响在线写入性能。