F1 如何用 AWS 智能体式 AI 将数据接入从 8 周压缩到 40 分钟

2026-08-04 59 预计阅读时间: 1 分钟
来源: aws.amazon.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:11 分钟

Formula 1® 与 AWS 合作构建 Data Accelerator,利用 Amazon Bedrock AgentCore 上的智能体式 AI 攭造其 MarTech 数据平台。最直观的变化是,接入一个新数据源不再需要最长 8 周,而是可以在大约 40 分钟内完成;与此同时,平台还开始自动处理 Schema 演进,并为车迷互动数据提供端到端可观测性。

这项改造的价值不只是“让 AI 写几段 SQL”。它把数据接入从一组依赖人工协调的任务,重构成了可规划、可执行、可验证、可追踪的工作流。

真正耗时的通常不是复制数据

MarTech 平台接入新数据源时,工程团队往往需要完成一条长链路:识别字段、判断数据类型、建立目标表、配置转换规则、检查隐私字段、部署任务、执行质量验证,再把运行状态接入监控系统。

单个步骤可能并不复杂,问题在于它们分散在工单、脚本、会议和不同团队之间。只要上游定义不清楚,流程就会退回到沟通阶段。所谓“8 周”,通常包含大量等待时间,而不是 8 周连续编码。

Data Accelerator 所体现的智能体式做法,是让系统围绕一个目标协调多步操作。例如,智能体可以读取数据源元数据,生成接入计划,调用平台工具创建资源,执行校验,并把每一步的输入、结果和异常写入统一的观测链路。人仍然负责策略、审批和异常决策,但不必手工推动每一个确定性步骤。

Schema 演进需要策略,不能只靠自动化

数据源上线后,字段仍会变化:营销工具可能增加活动属性,应用事件可能调整类型,合作伙伴也可能删除或重命名字段。如果平台只在首次接入时生成固定映射,后续变更很容易让任务失败,或者更危险地造成静默数据丢失。

自动化 Schema 演进应当先对变更分类:

  • 新增可选字段通常可以自动接受,并同步更新目录和目标表。
  • 类型从整数扩大为小数,可能可以执行兼容性升级。
  • 字段删除、重命名或类型收窄,应该暂停并请求人工审批。
  • 涉及邮箱、设备标识或其他敏感数据的字段,需要重新执行治理策略。

这也是智能体与普通脚本的差别之一:脚本通常执行固定分支,而智能体工作流可以结合元数据、组织策略和工具反馈,决定继续执行、修复重试,还是把问题交给工程师。不过,生产环境中的决策边界必须由明确规则约束,不能把高风险变更完全交给模型自行判断。

可以这样实践:先搭一个可审计的接入工作流

下面是一个可直接运行的本地示例。它不是 F1 或 AWS 的真实实现,而是根据摘要抽象出的最小工作流:读取当前 Schema 和候选 Schema,对变更分类,生成接入结果,并输出可供日志系统采集的结构化事件。

将以下内容保存为 data_accelerator_demo.py,使用 Python 3.10 或更高版本运行,无需安装第三方依赖:

from __future__ import annotations

import json
import sys
import time
import uuid
from pathlib import Path

SAFE_WIDENING = {
    ("integer", "number"),
}


def load_schema(path: str) -> dict[str, str]:
    data = json.loads(Path(path).read_text(encoding="utf-8"))
    if not isinstance(data, dict) or not all(
        isinstance(k, str) and isinstance(v, str) for k, v in data.items()
    ):
        raise ValueError("schema must be a JSON object of field: type pairs")
    return data


def classify_changes(current: dict[str, str], candidate: dict[str, str]) -> list[dict]:
    changes = []
    for field in sorted(current.keys() | candidate.keys()):
        old_type = current.get(field)
        new_type = candidate.get(field)

        if old_type is None:
            decision = "auto_apply"
            kind = "field_added"
        elif new_type is None:
            decision = "manual_review"
            kind = "field_removed"
        elif old_type == new_type:
            continue
        elif (old_type, new_type) in SAFE_WIDENING:
            decision = "auto_apply"
            kind = "type_widened"
        else:
            decision = "manual_review"
            kind = "type_changed"

        changes.append({
            "field": field,
            "change": kind,
            "from": old_type,
            "to": new_type,
            "decision": decision,
        })
    return changes


def emit(run_id: str, stage: str, status: str, **details: object) -> None:
    event = {
        "timestamp": int(time.time()),
        "run_id": run_id,
        "stage": stage,
        "status": status,
        "details": details,
    }
    print(json.dumps(event, ensure_ascii=False))


def main() -> int:
    if len(sys.argv) != 3:
        print("Usage: python data_accelerator_demo.py current.json candidate.json")
        return 2

    run_id = str(uuid.uuid4())
    emit(run_id, "schema_discovery", "started")

    try:
        current = load_schema(sys.argv[1])
        candidate = load_schema(sys.argv[2])
        emit(run_id, "schema_discovery", "completed", fields=len(candidate))

        changes = classify_changes(current, candidate)
        requires_review = any(c["decision"] == "manual_review" for c in changes)
        emit(
            run_id,
            "schema_evolution",
            "blocked" if requires_review else "approved",
            changes=changes,
        )

        result = {
            "run_id": run_id,
            "deployable": not requires_review,
            "changes": changes,
        }
        Path("onboarding-result.json").write_text(
            json.dumps(result, indent=2, ensure_ascii=False),
            encoding="utf-8",
        )
        return 1 if requires_review else 0
    except Exception as exc:
        emit(run_id, "workflow", "failed", error=str(exc))
        return 1


if __name__ == "__main__":
    raise SystemExit(main())

创建两份测试 Schema:

printf '%s\n' '{"fan_id":"string","sessions":"integer","email":"string"}' > current.json
printf '%s\n' '{"fan_id":"string","sessions":"number","email":"string","campaign_id":"string"}' > candidate.json
python data_accelerator_demo.py current.json candidate.json
cat onboarding-result.json

这个例子会把 sessions 从整数到数值的变化识别为安全扩展,并自动接受新增的 campaign_id。如果删除 email,或把它改成其他类型,工作流会返回非零退出码并要求人工检查。

接入 Amazon Bedrock AgentCore 时,可以将这里的 load_schemaclassify_changes、部署和质量检查分别封装为受控工具,由智能体负责编排。工具接口应使用结构化输入输出,并限制资源范围;数据库凭证、写入权限和治理规则不应直接放进提示词。

可观测性必须覆盖整条决策链

只记录“任务成功”不足以解释一个智能体工作流。端到端可观测性至少要回答:哪个请求触发了接入、智能体选择了哪些工具、每次调用修改了什么、Schema 为什么获批或被阻止、重试发生在哪里,以及最终数据质量是否达标。

实践中可以统一记录以下字段:

  • run_id:贯穿规划、工具调用、部署和验证的关联标识。
  • source_id 与 Schema 版本:确定本次运行处理了什么。
  • 决策和依据:区分规则决策、模型建议与人工审批。
  • 工具调用耗时、重试次数和错误类型:定位性能与稳定性问题。
  • 输入输出摘要:支持审计,同时避免把敏感原始数据写入日志。
  • 数据质量指标:包括行数差异、空值率、重复率和延迟。

模型调用成功不等于数据接入成功。真正的完成条件应包含目标资源已创建、转换任务已运行、质量门槛已通过、目录已更新,并且监控与告警已经生效。

落地时守住三个边界

从数周降到分钟很有吸引力,但团队不应只衡量平均接入时间。上线前还需要检查三个边界:

  1. 权限边界:智能体只能调用白名单工具,并使用最小权限身份执行操作。
  2. 审批边界:破坏性 Schema 变化、敏感字段和生产写入必须设置人工门禁。
  3. 验证边界:模型输出不能直接视为事实,所有资源配置和数据映射都要经过确定性校验。

更稳妥的采用路径,是先自动化元数据发现、计划生成和只读验证,再逐步开放低风险写操作。用接入耗时、人工触点数量、自动修复率、回滚率和质量事故数共同衡量效果,才能判断智能体究竟减少了运营负担,还是仅仅把复杂性转移到了新的系统里。


相关推荐