Boomi Scribe 如何用 AWS 把集成 DAG 自动变成可维护文档

2026-09-01 38 预计阅读时间: 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.

预计阅读时间:9 分钟

企业集成流程经常比业务代码更难维护:节点多、依赖复杂、版本变化频繁,而文档通常滞后于实际配置。Boomi Scribe 的思路是把集成工作流的有向无环图(DAG)交给 AI 代理处理,自动解析组件关系、生成详细说明,并在大规模场景下比较组件版本。

从已披露的技术栈看,这套能力由 Amazon Bedrock、Amazon SageMaker AI、Amazon S3、Amazon DynamoDB 和 AWS Lambda 共同支撑。关键并不是简单地“让模型读一段 JSON”,而是建立一条可追踪、可重试、可评估的文档生成流水线。

从 DAG 到文档,需要拆成多段处理

集成 DAG 通常包含节点、连接边、组件配置和版本信息。直接把完整 DAG 塞进提示词,会遇到几个实际问题:超出模型上下文、敏感配置泄露、输出结构不稳定,以及无法精确定位失败节点。

更稳妥的处理链可以分成四步:

  1. 接收与归档:把原始 DAG 和版本快照写入 Amazon S3,保留可重放的输入。
  2. 解析与标准化:由 AWS Lambda 提取节点、边、组件类型和版本,生成稳定的中间表示。
  3. 生成与评估:调用 Amazon Bedrock 中的基础模型生成文档;Amazon SageMaker AI 可以承担定制模型、评估或其他机器学习任务。
  4. 索引与状态管理:在 Amazon DynamoDB 中记录任务状态、输入版本、输出位置和错误信息。

这几个服务承担不同职责。S3 保存体积较大的原始 DAG 与文档,DynamoDB 保存适合按任务或组件查询的元数据,Lambda 负责事件驱动编排,Bedrock 负责生成式推理。这样的边界也让失败重试更明确:生成失败时,不必重新上传和解析整个工作流。

版本比较比生成摘要更有价值

一次性的流程说明很容易生成,真正困难的是回答“这次发布改了什么”。组件版本比较需要先对 DAG 做结构化归一化,再把差异交给模型解释。

例如,可以为每个节点构造稳定标识:

component_key = workflow_id + node_id + component_type

比较两个版本时,程序先计算确定性的变化:

  • 新增或删除了哪些节点;
  • 哪些连接边发生变化;
  • 组件版本、关键参数或错误处理策略是否变化;
  • 变化影响了哪些下游节点。

模型随后负责把结构化差异转成开发者能读懂的发布说明,而不是自行猜测两个大型 JSON 的差别。这能减少遗漏,也便于在测试中验证差异计算逻辑。

可以这样实践:用 Lambda 生成 DAG 文档

下面是一个可改造的最小示例,并非 Boomi Scribe 的原始实现。它假设 S3 的 incoming/ 目录收到 DAG JSON 后触发 Lambda,函数调用 Amazon Bedrock 生成 Markdown,再把结果写回 S3,并在 DynamoDB 中登记状态。

运行前需要设置 MODEL_IDTABLE_NAME 环境变量,并为 Lambda 授予 s3:GetObjects3:PutObjectbedrock:InvokeModeldynamodb:PutItem 权限。模型 ID 应替换为账户所在区域中已启用、支持 Bedrock Converse API 的模型。

import json
import os
import time
import urllib.parse

import boto3

s3 = boto3.client("s3")
bedrock = boto3.client("bedrock-runtime")
dynamodb = boto3.resource("dynamodb")

table = dynamodb.Table(os.environ["TABLE_NAME"])
model_id = os.environ["MODEL_ID"]


def normalize_dag(dag):
    nodes = sorted(
        [
            {
                "id": node["id"],
                "type": node.get("type", "unknown"),
                "version": node.get("version", "unknown"),
                "description": node.get("description", ""),
            }
            for node in dag.get("nodes", [])
        ],
        key=lambda node: node["id"],
    )
    edges = sorted(
        [
            {"from": edge["from"], "to": edge["to"]}
            for edge in dag.get("edges", [])
        ],
        key=lambda edge: (edge["from"], edge["to"]),
    )
    return {"workflow_id": dag["workflow_id"], "nodes": nodes, "edges": edges}


def lambda_handler(event, context):
    record = event["Records"][0]
    bucket = record["s3"]["bucket"]["name"]
    key = urllib.parse.unquote_plus(record["s3"]["object"]["key"])

    if not key.startswith("incoming/"):
        return {"statusCode": 200, "body": "Ignored non-input object"}

    dag = json.loads(s3.get_object(Bucket=bucket, Key=key)["Body"].read())
    normalized = normalize_dag(dag)
    workflow_id = normalized["workflow_id"]
    generated_at = int(time.time())

    prompt = """You are documenting an enterprise integration workflow.
Use only the supplied DAG. Do not invent systems, credentials, or behavior.
Return Markdown with these sections:
1. Purpose
2. Component inventory
3. Execution flow
4. Dependencies and failure points
5. Version notes

DAG:
""" + json.dumps(normalized, ensure_ascii=False)

    response = bedrock.converse(
        modelId=model_id,
        messages=[{"role": "user", "content": [{"text": prompt}]}],
        inferenceConfig={"maxTokens": 2000, "temperature": 0.1},
    )
    markdown = response["output"]["message"]["content"][0]["text"]
    output_key = f"documentation/{workflow_id}/{generated_at}.md"

    s3.put_object(
        Bucket=bucket,
        Key=output_key,
        Body=markdown.encode("utf-8"),
        ContentType="text/markdown; charset=utf-8",
    )
    table.put_item(
        Item={
            "workflow_id": workflow_id,
            "generated_at": generated_at,
            "source_key": key,
            "output_key": output_key,
            "status": "COMPLETED",
        }
    )

    return {"statusCode": 200, "body": json.dumps({"output_key": output_key})}

可用下面的测试 DAG 作为输入:

{
  "workflow_id": "order-sync",
  "nodes": [
    {"id": "read-orders", "type": "database-reader", "version": "2.1"},
    {"id": "map-fields", "type": "transform", "version": "4.0"},
    {"id": "send-crm", "type": "http-client", "version": "3.2"}
  ],
  "edges": [
    {"from": "read-orders", "to": "map-fields"},
    {"from": "map-fields", "to": "send-crm"}
  ]
}

生产实现还应把状态先写成 PROCESSING,捕获异常后更新为 FAILED,并配置死信队列或重试策略。对于大型 DAG,可以按子图生成文档片段,再用一次汇总调用拼装总览。

规模化时必须守住的边界

文档生成涉及企业连接器、字段映射和端点配置,输入清洗不能省略。把 DAG 发给模型前,应删除密码、令牌、连接字符串和不需要进入文档的业务数据。S3 对象、DynamoDB 表和日志也应采用最小权限与加密策略。

模型输出同样不能直接视为事实。可以在生成后增加几类自动校验:文档提到的组件 ID 必须存在于 DAG;组件数量必须匹配;所有版本号必须来自结构化输入;缺失节点或无效引用应阻止发布。

采用这类方案时,可以从一个低风险工作流开始,重点观察四项指标:生成成功率、人工修改比例、版本差异准确率和单次生成成本。只有当结构化解析与确定性校验稳定后,AI 生成的文档才适合进入正式发布流程。


相关推荐