Administrator
发布于 2026-07-12 / 10 阅读
1

DynamoDB 跨账号最小停机迁移 —— 演练完整记录

DynamoDB 跨账号最小停机迁移 —— 演练完整记录

演练日期:2026-06-22。目的:验证「S3 Export/Import + DynamoDB Streams + Lambda」跨账号最小停机迁移方案的可行性。
结论:全流程数据零丢失,源表与目标表精确对齐(1,562,087 == 1,562,087),内容抽样 50/50 一致,Lambda 11000+ 次调用 0 错误。


0. 总览

0.1 架构与数据流

                    源账号 source_account_id (cn-north-1)            目标账号 target_account_id (cn-north-1)
                    ┌──────────────────────────┐               ┌──────────────────────────┐
  模拟业务写入  ──► │ allinone_qa_LargeDataTable │               │ allinone_qa_LargeDataTable │
 (simulate_      │   ├─ PITR (全量基线来源)    │  Export S3    │   (Import from S3 新建)    │
  workload.py)   │   └─ Streams(增量来源)      │ ────────────► │                            │
                    └──────────┬───────────────┘   menglonq桶   └──────────▲───────────────┘
                               │ Stream                                     │ 跨账号写
                               ▼                                            │ (资源策略授权)
                    ┌──────────────────────────┐                          │
                    │ Lambda: ddb-xacct-        │ ─────────────────────────┘
                    │ replicator (事件源映射)    │
                    └──────────────────────────┘
  • 全量通道:Export to S3(基于 PITR)→ Import from S3 在目标账号建新表。
  • 增量通道:源表 Streams → 源账号 Lambda → 跨账号写目标表(方式 B:目标表资源策略授权,免 AssumeRole)。
  • 衔接保证:Streams 先于 Export 开启(制造重叠)+ Lambda 内 EXPORT_TS 过滤基线前记录 + 时间戳幂等条件写。

0.2 账号与身份

角色账号 ID接入方式实际身份 ARN
源账号source_account_iddefault profilearn:aws-cn:iam::source_account_id:user/menglonq
目标账号target_account_idEC2 实例角色(instance profile)arn:aws-cn:sts::target_account_id:assumed-role/AmazonSSMRoleForInstancesQuickSetup/i-0dc0c796dc5ba44ef

0.3 运行环境

  • 主机:EC2,Windows Server 2019,北京区
  • Python 3.12.2 / boto3 1.37.11 / botocore 1.37.11 / 8 核
  • 工作目录:C:\Users\Administrator\Documents\CaseHelper\migration_test

0.4 交付脚本清单

文件作用
load_data.py批量补写基线测试数据
simulate_workload.py模拟上游业务持续增删改
lambda/replicator.py跨账号增量复制 Lambda 函数体
deploy_lambda.py在源账号部署角色 + 函数 + 事件源映射
watch_iterator_age.pyCloudWatch IteratorAge 追平监控

1. 阶段零:复用表与补写基线数据

1.1 复用表的探测

复用已存在的 allinone_qa_LargeDataTable(原约 110w 条)。探测得到结构:

主键id(S, HASH) + timestamp(S, RANGE) —— 复合主键
计费模式PAY_PER_REQUEST(按需,无需预置容量)
StreamNone(建表时未开,符合要求)
初始 ItemCount(~6h 近似)1,100,000
TableSizeBytes2,462,853,540

数据 schema(每条 item 的属性)

属性类型说明
idSuuid4 字符串,分区键
timestampSISO8601,排序键,如 2026-03-13T12:47:30.385664
random_stringS超长随机串,约 2048 字符
categoryS单字符,取值 A/B/C/D/E
numeric_fieldN整数,约 1~1,000,000
metadataMMap:created_by(S) / version(S, v1v3) / tags(L, 510 个 tag_n)

1.2 load_data.py 设计与参数讲究

"""
向 DynamoDB 表批量补写测试数据,schema 复刻表内现有"基础数据"。

用途:为 DynamoDB 跨账号迁移方案搭建测试基线数据。
表:allinone_qa_LargeDataTable(cn-north-1),主键 id(HASH,S) + timestamp(RANGE,S)。

显式锁定 default profile + cn-north-1,避免被环境变量 AWS_PROFILE=global / AWS_REGION 干扰。

用法:
    python load_data.py                # 默认补写 400000 条
    python load_data.py --count 100000 # 自定义条数
    python load_data.py --workers 16
"""

import argparse
import random
import string
import threading
import time
import uuid
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timezone

import boto3
from botocore.config import Config
from botocore.exceptions import ClientError

# ----------------------- 固定配置 -----------------------
PROFILE = "default"
REGION = "cn-north-1"
TABLE = "allinone_qa_LargeDataTable"

BATCH_SIZE = 25  # BatchWriteItem 单批上限
CATEGORIES = ["A", "B", "C", "D", "E"]
VERSIONS = ["v1", "v2", "v3"]
# 现有数据 random_string 长度约 2000+;用区间随机以贴近真实分布
RAND_MIN, RAND_MAX = 1800, 2200

_ALPHABET = string.ascii_letters + string.digits

# 进度计数(跨线程)
_lock = threading.Lock()
_written = 0
_failed = 0


def _make_client():
    """构造写表用的 DynamoDB client,连接池放大以支撑多线程并发。"""
    session = boto3.Session(profile_name=PROFILE, region_name=REGION)
    cfg = Config(
        max_pool_connections=64,
        retries={"max_attempts": 10, "mode": "adaptive"},
    )
    return session.client("dynamodb", config=cfg)


def _rand_string(n):
    return "".join(random.choices(_ALPHABET, k=n))


def _gen_item():
    """生成一条与现有基础数据同 schema 的 item(DynamoDB JSON 低层格式)。"""
    tag_count = random.randint(5, 10)
    return {
        "id": {"S": str(uuid.uuid4())},
        "timestamp": {"S": datetime.now(timezone.utc).isoformat()},
        "random_string": {"S": _rand_string(random.randint(RAND_MIN, RAND_MAX))},
        "category": {"S": random.choice(CATEGORIES)},
        "numeric_field": {"N": str(random.randint(1, 1_000_000))},
        "metadata": {
            "M": {
                "created_by": {"S": "load_script"},
                "version": {"S": random.choice(VERSIONS)},
                "tags": {"L": [{"S": f"tag_{i}"} for i in range(tag_count)]},
            }
        },
    }


def _write_batch(client, items):
    """写一批(≤25 条),对 UnprocessedItems 做指数退避重试。返回成功条数。"""
    request = {TABLE: [{"PutRequest": {"Item": it}} for it in items]}
    attempt = 0
    while request.get(TABLE):
        try:
            resp = client.batch_write_item(RequestItems=request)
        except ClientError as e:
            # 限流等可重试错误:退避后整批重试(已配 adaptive retries 兜底)
            attempt += 1
            if attempt > 10:
                raise
            time.sleep(min(2 ** attempt * 0.05, 5))
            continue
        unprocessed = resp.get("UnprocessedItems", {})
        if not unprocessed.get(TABLE):
            break
        request = unprocessed
        attempt += 1
        time.sleep(min(2 ** attempt * 0.05, 5))
    return len(items)


def _worker(client, n_items, report_every):
    """单线程任务:生成并写入 n_items 条。"""
    global _written, _failed
    buf = []
    done_local = 0
    for _ in range(n_items):
        buf.append(_gen_item())
        if len(buf) == BATCH_SIZE:
            done_local += _flush(client, buf)
            buf = []
            if done_local >= report_every:
                _bump(done_local, 0)
                done_local = 0
    if buf:
        done_local += _flush(client, buf)
    if done_local:
        _bump(done_local, 0)


def _flush(client, buf):
    global _failed
    try:
        return _write_batch(client, buf)
    except Exception as e:  # noqa: BLE001 - 单批失败不应中断整体载入
        with _lock:
            _failed += len(buf)
        print(f"[WARN] batch failed ({len(buf)} items): {e}")
        return 0


def _bump(ok, bad):
    global _written, _failed
    with _lock:
        _written += ok
        _failed += bad


def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--count", type=int, default=400_000, help="补写条数")
    parser.add_argument("--workers", type=int, default=16, help="并发线程数")
    args = parser.parse_args()

    total = args.count
    workers = args.workers
    client = _make_client()

    # 平均分配到各线程
    per = total // workers
    remainder = total % workers
    chunks = [per + (1 if i < remainder else 0) for i in range(workers)]

    print(f"目标表    : {TABLE} ({REGION}, profile={PROFILE})")
    print(f"补写条数  : {total:,}  并发线程: {workers}  每批: {BATCH_SIZE}")
    print(f"线程分配  : {chunks}")
    print("-" * 60)

    start = time.time()

    # 后台进度打印
    stop_flag = {"stop": False}

    def _progress():
        while not stop_flag["stop"]:
            time.sleep(5)
            with _lock:
                w, f = _written, _failed
            elapsed = time.time() - start
            rate = w / elapsed if elapsed > 0 else 0
            pct = w / total * 100 if total else 100
            eta = (total - w) / rate if rate > 0 else 0
            print(f"  进度 {w:>9,}/{total:,} ({pct:5.1f}%)  "
                  f"失败 {f}  速率 {rate:8.0f} 条/s  ETA {eta:6.0f}s")

    monitor = threading.Thread(target=_progress, daemon=True)
    monitor.start()

    report_every = max(1000, total // workers // 20)
    with ThreadPoolExecutor(max_workers=workers) as ex:
        futures = [ex.submit(_worker, client, c, report_every) for c in chunks]
        for fut in as_completed(futures):
            fut.result()  # 抛出未捕获异常(理论上 _worker 内已吞批级错误)

    stop_flag["stop"] = True
    elapsed = time.time() - start
    print("-" * 60)
    print(f"完成:成功 {_written:,} 条,失败 {_failed} 条,耗时 {elapsed:.1f}s,"
          f"平均 {_written/elapsed if elapsed else 0:.0f} 条/s")


if __name__ == "__main__":
    main()

目标:补写 40w 条,schema 完全复刻现有数据。关键设计:

PROFILE = "default"; REGION = "cn-north-1"; TABLE = "allinone_qa_LargeDataTable"
BATCH_SIZE = 25                    # BatchWriteItem 单批硬上限就是 25
RAND_MIN, RAND_MAX = 1800, 2200    # random_string 长度区间,贴近现有数据的 ~2048

参数讲究

  1. BatchWriteItem 批量 25 条:DynamoDB BatchWriteItem 单次最多 25 个写请求,这是 API 硬上限,不是调优值。
  2. 连接池放大Config(max_pool_connections=64),否则多线程会因默认连接池(10)排队。
  3. 自适应重试retries={"max_attempts": 10, "mode": "adaptive"},按需表初期可能触发 ThrottlingException,adaptive 模式会自动退避。
  4. UnprocessedItems 重试BatchWriteItem 即使 200 也可能部分未写(返回在 UnprocessedItems),必须循环重试这部分,否则静默丢数据:
def _write_batch(client, items):
    request = {TABLE: [{"PutRequest": {"Item": it}} for it in items]}
    attempt = 0
    while request.get(TABLE):
        resp = client.batch_write_item(RequestItems=request)
        unprocessed = resp.get("UnprocessedItems", {})
        if not unprocessed.get(TABLE):
            break
        request = unprocessed
        attempt += 1
        time.sleep(min(2 ** attempt * 0.05, 5))   # 指数退避
    return len(items)
  1. 16 线程并发:8 核机器,I/O 密集(网络写),16 线程吞吐最佳。线程分配用 per + (1 if i < remainder else 0) 均摊余数。

1.3 执行与结果

# 先冒烟 100 条
python load_data.py --count 100 --workers 4      # 100 条 0 失败,132 条/s

# 全量 40w(后台运行,PYTHONIOENCODING=utf-8 避免 Windows GBK 控制台中文乱码)
PYTHONIOENCODING=utf-8 python load_data.py --count 400000 --workers 16

结果:成功 400,000 条,失败 0 条,耗时 188.1s,平均 2127 条/s。表内基础数据补到约 150w。

注:Python print 重定向到文件时是块缓冲,后台运行期间输出文件长时间为 0 字节是正常的,进程结束才 flush。判断进度可读输出文件或 ps -W 看进程存活。


2. 阶段一:模拟业务客户端 simulate_workload.py

"""
模拟真实业务应用对 DynamoDB 表的持续增删改(INSERT / UPDATE / DELETE)。

用途:DynamoDB 跨账号迁移方案验证。在执行 Export -> Import -> Streams+Lambda
增量同步的全过程中,让源表持续产生变更,以检验:
  - 全量基线导出期间的写入是否被增量通道捕获
  - cutover 前增量能否追平、数据是否最终一致

设计:
  - 维护一个进程内"已知主键池"(known keys):启动时从表抽样捞一批,运行中新 INSERT 的
    key 也加入池;UPDATE/DELETE 从池中随机取,保证操作命中真实存在的 item。
  - 三类操作按 --insert/--update/--delete 权重随机混合。
  - 通过 --tps 控制目标吞吐(令牌桶节流),--duration 控制运行时长,--workers 控制并发。
  - 每条变更都在 metadata.last_op / metadata.op_ts 留痕,便于迁移后做数据比对核验。

显式锁定 default profile + cn-north-1,避免被环境变量 AWS_PROFILE=global 干扰。

用法:
    python simulate_workload.py --tps 50 --duration 600
    python simulate_workload.py --tps 100 --duration 0   # duration=0 表示一直跑到 Ctrl+C
    python simulate_workload.py --insert 3 --update 6 --delete 1 --tps 80
"""

import argparse
import random
import signal
import string
import threading
import time
import uuid
from collections import deque
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime, timezone

import boto3
from botocore.config import Config
from botocore.exceptions import ClientError

# ----------------------- 固定配置 -----------------------
PROFILE = "default"
REGION = "cn-north-1"
TABLE = "allinone_qa_LargeDataTable"

CATEGORIES = ["A", "B", "C", "D", "E"]
VERSIONS = ["v1", "v2", "v3"]
RAND_MIN, RAND_MAX = 1800, 2200
_ALPHABET = string.ascii_letters + string.digits

# 启动时从表抽样捞多少 key 进池(供 UPDATE/DELETE 命中已有数据)
SEED_KEY_TARGET = 5000

_stop = threading.Event()


# ----------------------- 主键池 -----------------------
class KeyPool:
    """线程安全的主键池,存 (id, timestamp) 元组。

    用 deque 限制上限,避免长时间运行内存膨胀;满了丢最旧的。
    """

    def __init__(self, maxlen=200_000):
        self._dq = deque(maxlen=maxlen)
        self._lock = threading.Lock()

    def add(self, key):
        with self._lock:
            self._dq.append(key)

    def random(self):
        with self._lock:
            if not self._dq:
                return None
            return random.choice(self._dq)

    def discard(self, key):
        # deque 无高效随机删除;DELETE 后不强制移除,靠条件写幂等兜底即可。
        # 简化处理:不移除,重复 DELETE 同 key 时第二次条件失败计为 noop。
        pass

    def __len__(self):
        with self._lock:
            return len(self._dq)


# ----------------------- 统计 -----------------------
class Stats:
    def __init__(self):
        self.lock = threading.Lock()
        self.insert = 0
        self.update = 0
        self.delete = 0
        self.noop = 0      # 条件写未命中(如 DELETE 已删的 key)
        self.error = 0

    def bump(self, field, n=1):
        with self.lock:
            setattr(self, field, getattr(self, field) + n)

    def snapshot(self):
        with self.lock:
            return (self.insert, self.update, self.delete, self.noop, self.error)


def _make_client():
    session = boto3.Session(profile_name=PROFILE, region_name=REGION)
    cfg = Config(max_pool_connections=64,
                 retries={"max_attempts": 10, "mode": "adaptive"})
    return session.client("dynamodb", config=cfg)


def _rand_string(n):
    return "".join(random.choices(_ALPHABET, k=n))


def _now_iso():
    return datetime.now(timezone.utc).isoformat()


def _seed_pool(client, pool, target):
    """启动时分页 scan 抽样,捞 target 个 key 进池。只取主键属性,省流量。"""
    scanned = 0
    kwargs = {
        "TableName": TABLE,
        "ProjectionExpression": "id, #ts",
        "ExpressionAttributeNames": {"#ts": "timestamp"},
        "Limit": 1000,
    }
    while len(pool) < target:
        resp = client.scan(**kwargs)
        for it in resp.get("Items", []):
            pool.add((it["id"]["S"], it["timestamp"]["S"]))
        scanned += resp.get("ScannedCount", 0)
        lek = resp.get("LastEvaluatedKey")
        if not lek:
            break
        kwargs["ExclusiveStartKey"] = lek
    print(f"[seed] 已抽样 {len(pool)} 个主键进池(scanned={scanned})")


# ----------------------- 三类操作 -----------------------
def _do_insert(client, pool, stats):
    key_id = str(uuid.uuid4())
    ts = _now_iso()
    item = {
        "id": {"S": key_id},
        "timestamp": {"S": ts},
        "random_string": {"S": _rand_string(random.randint(RAND_MIN, RAND_MAX))},
        "category": {"S": random.choice(CATEGORIES)},
        "numeric_field": {"N": str(random.randint(1, 1_000_000))},
        "metadata": {"M": {
            "created_by": {"S": "workload_sim"},
            "version": {"S": random.choice(VERSIONS)},
            "last_op": {"S": "INSERT"},
            "op_ts": {"S": ts},
            "tags": {"L": [{"S": f"tag_{i}"} for i in range(random.randint(5, 10))]},
        }},
    }
    client.put_item(TableName=TABLE, Item=item)
    pool.add((key_id, ts))
    stats.bump("insert")


def _do_update(client, pool, stats):
    key = pool.random()
    if key is None:
        stats.bump("noop")
        return
    key_id, ts = key
    now = _now_iso()
    try:
        client.update_item(
            TableName=TABLE,
            Key={"id": {"S": key_id}, "timestamp": {"S": ts}},
            UpdateExpression=("SET numeric_field = :n, category = :c, "
                              "metadata.last_op = :op, metadata.op_ts = :t"),
            ExpressionAttributeValues={
                ":n": {"N": str(random.randint(1, 1_000_000))},
                ":c": {"S": random.choice(CATEGORIES)},
                ":op": {"S": "UPDATE"},
                ":t": {"S": now},
            },
            # 仅当该 item 仍存在时才更新(已被 DELETE 的不重建)
            ConditionExpression="attribute_exists(id)",
        )
        stats.bump("update")
    except ClientError as e:
        if e.response["Error"]["Code"] == "ConditionalCheckFailedException":
            stats.bump("noop")
        else:
            raise


def _do_delete(client, pool, stats):
    key = pool.random()
    if key is None:
        stats.bump("noop")
        return
    key_id, ts = key
    try:
        client.delete_item(
            TableName=TABLE,
            Key={"id": {"S": key_id}, "timestamp": {"S": ts}},
            ConditionExpression="attribute_exists(id)",
        )
        stats.bump("delete")
    except ClientError as e:
        if e.response["Error"]["Code"] == "ConditionalCheckFailedException":
            stats.bump("noop")  # 已被删过,幂等 noop
        else:
            raise


def _pick_op(weights):
    r = random.random() * weights["total"]
    if r < weights["insert"]:
        return _do_insert
    if r < weights["insert"] + weights["update"]:
        return _do_update
    return _do_delete


# ----------------------- 令牌桶节流 -----------------------
class RateLimiter:
    """简单令牌桶:每秒补充 tps 个令牌。tps<=0 表示不限速。"""

    def __init__(self, tps):
        self.tps = tps
        self.lock = threading.Lock()
        self.allowance = float(tps)
        self.last = time.monotonic()

    def acquire(self):
        if self.tps <= 0:
            return
        while True:
            with self.lock:
                now = time.monotonic()
                self.allowance += (now - self.last) * self.tps
                self.last = now
                if self.allowance > self.tps:
                    self.allowance = self.tps
                if self.allowance >= 1:
                    self.allowance -= 1
                    return
            time.sleep(0.002)


def _worker(client, pool, stats, weights, limiter):
    while not _stop.is_set():
        limiter.acquire()
        if _stop.is_set():
            break
        op = _pick_op(weights)
        try:
            op(client, pool, stats)
        except Exception as e:  # noqa: BLE001 - 单次操作失败不中断模拟
            stats.bump("error")
            print(f"[ERR] {op.__name__}: {e}")


def main():
    p = argparse.ArgumentParser()
    p.add_argument("--tps", type=int, default=50, help="目标总吞吐(操作/秒),<=0 不限速")
    p.add_argument("--duration", type=int, default=300, help="运行秒数,0=直到 Ctrl+C")
    p.add_argument("--workers", type=int, default=8, help="并发线程数")
    p.add_argument("--insert", type=int, default=3, help="INSERT 权重")
    p.add_argument("--update", type=int, default=5, help="UPDATE 权重")
    p.add_argument("--delete", type=int, default=2, help="DELETE 权重")
    p.add_argument("--seed-keys", type=int, default=SEED_KEY_TARGET,
                   help="启动时抽样进池的 key 数")
    args = p.parse_args()

    weights = {
        "insert": args.insert,
        "update": args.update,
        "delete": args.delete,
        "total": args.insert + args.update + args.delete,
    }

    client = _make_client()
    pool = KeyPool()

    print(f"目标表    : {TABLE} ({REGION}, profile={PROFILE})")
    print(f"操作权重  : INSERT={args.insert} UPDATE={args.update} DELETE={args.delete}")
    print(f"目标 TPS  : {args.tps}  时长: {args.duration or '∞'}s  线程: {args.workers}")
    print("-" * 64)

    _seed_pool(client, pool, args.seed_keys)
    if len(pool) == 0:
        print("[警告] 主键池为空,UPDATE/DELETE 将全部 noop,仅 INSERT 生效。")

    # Ctrl+C 优雅停止
    signal.signal(signal.SIGINT, lambda *_: _stop.set())

    limiter = RateLimiter(args.tps)
    stats = Stats()
    start = time.time()

    def _progress():
        last = stats.snapshot()
        last_t = start
        while not _stop.is_set():
            time.sleep(5)
            now = time.time()
            cur = stats.snapshot()
            di = sum(cur[:3]) - sum(last[:3])
            cur_tps = di / (now - last_t) if now > last_t else 0
            i, u, d, n, e = cur
            print(f"  t+{now-start:5.0f}s | INS {i:>7,} UPD {u:>7,} DEL {d:>7,} "
                  f"| noop {n:>6,} err {e} | 池 {len(pool):>6,} | ~{cur_tps:5.0f} ops/s")
            last, last_t = cur, now

    threading.Thread(target=_progress, daemon=True).start()

    with ThreadPoolExecutor(max_workers=args.workers) as ex:
        for _ in range(args.workers):
            ex.submit(_worker, client, pool, stats, weights, limiter)
        # 主线程负责计时
        try:
            while not _stop.is_set():
                if args.duration and (time.time() - start) >= args.duration:
                    _stop.set()
                    break
                time.sleep(0.2)
        except KeyboardInterrupt:
            _stop.set()

    i, u, d, n, e = stats.snapshot()
    elapsed = time.time() - start
    total_ops = i + u + d
    print("-" * 64)
    print(f"运行结束:耗时 {elapsed:.1f}s")
    print(f"  INSERT={i:,}  UPDATE={u:,}  DELETE={d:,}  noop={n:,}  error={e}")
    print(f"  有效变更 {total_ops:,} 次,平均 {total_ops/elapsed if elapsed else 0:.1f} ops/s")
    print(f"  净增条数 ≈ INSERT - DELETE = {i - d:,}")


if __name__ == "__main__":
    main()

模拟"上游业务无法停机"的真实写入流,持续产生 INSERT/UPDATE/DELETE,喂给后续增量通道。

2.1 三个核心组件

(1) KeyPool —— 已知主键池(线程安全)

UPDATE/DELETE 必须命中真实存在的 item,否则全是 noop。设计:

  • 启动时 scan 抽样捞一批主键进池(默认 5000,演练用 3000)。
  • 运行中每次 INSERT 的新 key 也加入池。
  • deque(maxlen=200000) 限上限,满了丢最旧,避免长跑内存膨胀。
  • DELETE 后不强制移除 key(deque 无高效随机删除),靠条件写幂等兜底——重复 DELETE 同 key 第二次条件失败计为 noop。

(2) RateLimiter —— 令牌桶限速

self.allowance += (now - self.last) * self.tps   # 按时间补充令牌
if self.allowance >= 1: self.allowance -= 1; return

每秒补充 tps 个令牌,控制目标吞吐。tps<=0 不限速。

(3) Stats —— 线程安全统计:insert / update / delete / noop / error 计数。

2.2 三类操作的实现讲究

操作实现幂等/防护
INSERTput_item 新 uuid,写 metadata.last_op=INSERTop_ts新 key 入池
UPDATEupdate_itemnumeric_field/category/metadata.last_op/op_tsConditionExpression="attribute_exists(id)",命中已删的 key → ConditionalCheckFailedException → 计 noop
DELETEdelete_item同上条件,已删过 → noop

留痕设计:每条变更写 metadata.last_op(操作类型)和 metadata.op_ts(操作时刻 ISO),用于迁移后精确比对最后的更新是否同步到位

操作权重默认 INSERT=3 / UPDATE=5 / DELETE=2(更新为主,贴近真实 OLTP)。

2.3 验证运行(60s)

PYTHONIOENCODING=utf-8 python simulate_workload.py --tps 40 --duration 60 --workers 8 --seed-keys 3000

结果

[seed] 已抽样 3289 个主键进池
INSERT=735  UPDATE=1,150  DELETE=447  noop=97  error=0
有效变更 2,332 次,平均 38.7 ops/s,净增 288 条
  • error=0:三类操作全部正常命中。
  • noop=97:预期内幂等行为(UPDATE/DELETE 命中已删 key),证明条件写生效。

正式演练时:模拟客户端在 Export 之前就启动(20:03:42,PID 见 §9),--duration 0 常驻,持续产生增量,直到 cutover 时手动停止。


3. 阶段二:开启 PITR 与 Streams(关键时序)

时序铁律:Streams 必须先于 Export 开启。 操作通过控制台完成,脚本核实结果:

实测值
StreamEnabledTrue
StreamViewTypeNEW_AND_OLD_IMAGES(增量回放 + 幂等判断都需要 new image)
LatestStreamArnarn:aws-cn:dynamodb:cn-north-1:source_account_id:table/allinone_qa_LargeDataTable/stream/2026-06-22T12:04:28.095
Stream 开启时刻2026-06-22T12:04:28.095 UTC = 北京 20:04:28
PITR StatusENABLED(Export to S3 的硬性前提)

为何先开 Stream:Stream 只能捕获开启之后的变更。先开 Stream(20:04:28),Export 在之后(20:08:47),二者间约 4 分钟重叠区间——这段变更同时进了基线快照和 Stream,靠幂等去重,不丢数据。若反过来(先导出后开流),则 [导出, 开流] 之间的变更两头都不在,永久丢失。


4. 阶段三:全量导出到 S3

通过控制台发起 Export,时间点选"当前时间"。脚本核实导出任务:

实测值
ExportStatusCOMPLETED
ExportTime2026-06-22 20:08:47.268 +08:00 = epoch 1782130127
ItemCount1,504,768
BilledSizeBytes2,462,853,540
S3Bucketmenglonq
S3BucketOwnertarget_account_id(目标账号自有桶)
S3Prefixtmp/ddb/backup/

关键决策:导出到目标账号的桶menglonq 桶 owner 是 target_account_id(目标账号),所以目标账号 Import 时读自己的桶,无需任何跨账号 S3 授权,最省事。

S3 产物结构(实测):

tmp/ddb/backup/AWSDynamoDB/01782130127268-c860c940/
    ├── manifest-summary.json
    ├── manifest-files.json
    └── data/                          # 24 个 .json.gz 文件,每个约 100 MB
        ├── 3goqtddyim5ynm75oghnyy6s7u.json.gz   (101,879,132 bytes)
        ├── 4aouzthtju5wpcsb2luvcnsgk4.json.gz   (101,466,287 bytes)
        └── ... (共 24 个)

抽样解压第一个 data 文件确认内容正常:单文件 64,257 行(item),属性结构与源表一致(random_string len=2048,metadata 为 Map 等)。

Export 异步、不消耗源表 RCU,导出期间模拟客户端持续写入,源表性能无影响。


5. 阶段四:跨账号导入(含 GZIP 坑)

在目标账号 target_account_id 通过 Import from S3 建新表。

5.1 ⚠️ 第一次导入失败:压缩类型选错

第一次导入(控制台未正确设置压缩类型):

ImportStatusFAILED
FailureCodeItemValidationError
ProcessedItemCount24(= data 目录文件数!)
ImportedItemCount0

根因诊断ProcessedItemCount=24 恰好等于 data/ 下的 24 个 .json.gz 文件数。这个"处理条数 = 文件数"的特征,说明导入把每个 gz 文件整体当成 1 个 item 去解析——即压缩类型选成了 NONE,拿 gzip 二进制当明文解析,每个文件报一个校验错。

诊断时还核实了:

  • 建表 KeySchema 正确(id HASH + timestamp RANGE,均 S)——排除主键问题。
  • InputFormat = DYNAMODB_JSON 正确。
  • S3 路径、BucketOwner 正确。
  • CloudWatchLogGroupArn: None——该导入未关联日志组,所以"check CloudWatch logs"无处可查(也是个坑:没配日志组时排错只能靠 describe_import 的特征)。

5.2 第二次导入成功:指定 GZIP

重新发起,压缩类型显式选 GZIP

ImportId01782130969812-476a6576
InputCompressionTypeGZIP
InputFormatDYNAMODB_JSON
S3KeyPrefixtmp/ddb/backup/AWSDynamoDB/01782130127268-c860c940/data/(指向 data 目录)
ImportStatusCOMPLETED
ProcessedItemCount1,504,768
ImportedItemCount1,504,768(== 导出 itemCount,零失败)✅

导入期间目标表 CREATING 不可读写;完成后 ACTIVE。导入建的新表不带 Stream(符合预期)。

结论:导出产物是 .json.gz,Import 的压缩类型必须选 GZIP。判断特征:处理条数 = 文件数、导入数 0、ItemValidationError


6. 阶段五:源账号部署增量复制 Lambda

在源账号 source_account_id 部署,由 deploy_lambda.py 一次建成三件套。部署前探测确认:源账号用户有 IAM/Lambda 权限、无同名资源、源表 Stream ARN 可取、EXPORT_TS 算定。

6.1 资源 1:IAM 执行角色 ddb-xacct-replicator-role

信任策略(允许 Lambda 服务 assume):

{"Version":"2012-10-17","Statement":[{"Effect":"Allow",
  "Principal":{"Service":"lambda.amazonaws.com"},"Action":"sts:AssumeRole"}]}

内联策略 replicator-inline(三段):

{
  "Version":"2012-10-17",
  "Statement":[
    {"Sid":"Logs","Effect":"Allow",
     "Action":["logs:CreateLogGroup","logs:CreateLogStream","logs:PutLogEvents"],
     "Resource":"arn:aws-cn:logs:*:*:*"},
    {"Sid":"ReadSourceStream","Effect":"Allow",
     "Action":["dynamodb:GetRecords","dynamodb:GetShardIterator",
               "dynamodb:DescribeStream","dynamodb:ListStreams"],
     "Resource":"arn:aws-cn:dynamodb:cn-north-1:source_account_id:table/allinone_qa_LargeDataTable/stream/*"},
    {"Sid":"WriteTargetTableCrossAccount","Effect":"Allow",
     "Action":["dynamodb:PutItem","dynamodb:UpdateItem","dynamodb:DeleteItem","dynamodb:BatchWriteItem"],
     "Resource":"arn:aws-cn:dynamodb:cn-north-1:target_account_id:table/allinone_qa_LargeDataTable"}
  ]
}

讲究

  • Stream 读权限 Resource 用 /stream/* 通配,因为 Stream ARN 带时间戳,重开 Stream 会变,通配避免硬编码失效。
  • 第三段是跨账号写目标表 ARN——这是「方式 B」中 identity 侧的一半授权(resource 侧的另一半在 §7)。
  • 新建角色后 sleep 10 等 IAM 传播,否则 create_function 偶发 InvalidParameterValue: role cannot be assumed

6.2 资源 2:Lambda 函数 ddb-xacct-replicator

配置讲究
Runtimepython3.12与本机一致
Handlerreplicator.lambda_handler
Timeout60单批 100 条几百毫秒就够,60s 极宽裕(15min 硬限是单批处理时间,非追平总时长)
MemorySize256 MBitem 含 2KB 串,256MB 足够
环境变量见下

环境变量

TARGET_TABLE_ARN = arn:aws-cn:dynamodb:cn-north-1:target_account_id:table/allinone_qa_LargeDataTable
TARGET_REGION    = cn-north-1
EXPORT_TS        = 1782130007
REPL_TS_ATTR     = _replication_ts

EXPORT_TS=1782130007 的算法:exportTime epoch 1782130127(北京 20:08:47)减 120 秒 = 1782130007(北京 20:06:47)。提前 2 分钟留 buffer,过滤掉确定已含在全量基线里的旧记录。这个提前量仍落在 Stream 开启时刻(20:04:28)之后,重叠窗口完整,不丢数据。

函数核心逻辑lambda/replicator.py):

  1. 方式 B 跨账号写:直接用 Lambda 自身执行角色凭证 + 目标表 ARN 写,不 AssumeRole
    _target = boto3.client("dynamodb", region_name=TARGET_REGION,
                           config=Config(retries={"max_attempts":10,"mode":"adaptive"}))
    # 写时 TableName=TARGET_TABLE_ARN(用完整 ARN 而非表名,触发资源策略鉴权)
    
  2. 基线过滤event_ts = Decimal(ddb["ApproximateCreationDateTime"]),若 < EXPORT_TS 直接跳过。
  3. 时间戳幂等条件写
    • INSERT/MODIFY:put_item(NewImage + _replication_ts=event_ts)
      ConditionExpression="attribute_not_exists(#ts) OR #ts < :ts"<,旧的覆盖不了新的)
    • REMOVE:delete_item(Keys)ConditionExpression="attribute_not_exists(#ts) OR #ts <= :ts"<=,防误删更新的写入)
  4. 错误分类ConditionalCheckFailedException 视为"幂等跳过"(INFO 日志,不计错误);其他异常加入 batchItemFailures 精确重试。
  5. NewImage 直接复用:Stream 的 NewImage 已是 DynamoDB JSON 低层格式,可直接当 put_item 的 Item。

6.3 资源 3:事件源映射(先禁用)

lam.create_event_source_mapping(
    EventSourceArn=SRC_STREAM_ARN,
    FunctionName="ddb-xacct-replicator",
    Enabled=False,                          # 先禁用!
    StartingPosition="TRIM_HORIZON",
    BatchSize=100,
    MaximumRetryAttempts=-1,
    BisectBatchOnFunctionError=True,
    FunctionResponseTypes=["ReportBatchItemFailures"],
)

得到 UUID:5706f37f-bcb7-43c0-85c9-6696cb5fbee7

每个参数的讲究

参数讲究
EnabledFalse目标表还在 CREATING 不可写、资源策略未配,先禁用,待就绪再 enable
StartingPositionTRIM_HORIZONStream 不支持 AT_TIMESTAMP,只能 TRIM_HORIZON(全重放)或 LATEST。用 TRIM_HORIZON 全重放 + 代码 EXPORT_TS 过滤,等效"从导出点消费"
BatchSize100单次最多攒 100 条调一次(上限 1000)。我们 item ~2KB,100 条约 200KB,远低于 6MB payload 上限
MaximumBatchingWindowInSeconds未设(=0)追平阶段要低延迟,有记录就立刻触发,不刻意攒批
FunctionResponseTypesReportBatchItemFailures函数返回精确失败记录,只重投失败那几条,不整批回退
BisectBatchOnFunctionErrorTrue整批异常时二分批重试,快速隔离毒丸记录
MaximumRetryAttempts-1(无限)⚠️ 见 §12 待改进:生产应配有限次 + DLQ,否则毒丸会卡死分片 24h

6.4 部署执行结果

[role] created ddb-xacct-replicator-role
[role] inline policy attached
[fn]   created ddb-xacct-replicator
[esm]  created 5706f37f-bcb7-43c0-85c9-6696cb5fbee7 (Enabled=False)

7. 阶段六:目标表挂资源策略(方式 B 的另一半)

待目标表 ACTIVE 后,在目标账号 target_account_id 给目标表挂资源策略。

前置确认:导入 COMPLETED、Imported=1,504,768、表 ACTIVE
风险点验证:原担心实例角色(AmazonSSMRoleForInstancesQuickSetup)无 dynamodb:PutResourcePolicy 权限——实测有权限,直接挂成功。

资源策略内容

{
  "Version":"2012-10-17",
  "Statement":[{
    "Sid":"AllowCrossAccountReplicatorWrites",
    "Effect":"Allow",
    "Principal":{"AWS":"arn:aws-cn:iam::source_account_id:role/ddb-xacct-replicator-role"},
    "Action":["dynamodb:PutItem","dynamodb:UpdateItem","dynamodb:DeleteItem","dynamodb:BatchWriteItem"],
    "Resource":"arn:aws-cn:dynamodb:cn-north-1:target_account_id:table/allinone_qa_LargeDataTable"
  }]
}

put_resource_policy 成功,RevisionId: 1782132298690,回读验证内容正确。

方式 B 授权完整链路(两侧都放行才写得通)

  • 源账号 identity 侧:执行角色内联策略允许写目标表 ARN(§6.1 第三段)
  • 目标账号 resource 侧:目标表资源策略授权该角色(本节)

8. 阶段七:启用映射与追平观察

8.1 启用事件源映射

lam.update_event_source_mapping(UUID="5706f37f-...", Enabled=True)
# State 秒切 Enabled(小表),LastProcessingResult: No records processed → 随即 OK

8.2 端到端验证(Lambda 确实在跨账号写)

enable 后约 2 分钟,目标表扫到 76 条 created_by=workload_sim 的记录,全部带 _replication_ts(基线导入的老数据无此属性)→ 证明是 Lambda 写入的增量。样本 last_op=UPDATE_replication_ts=1782131346,连增量 UPDATE 都正确同步。

8.3 IteratorAge 追平曲线(CloudWatch AWS/Lambda,Maximum)

时刻      IteratorAge(Max)     说明
20:55      3037.9 s          ┐ 阶段2:积压下降中(模拟客户端先写了~50分钟才enable,
20:56      2769.1 s          │         TRIM_HORIZON 从最老记录读起,故起跳≈积压时长)
20:57      2406.2 s          │
20:58      1706.8 s          │
20:59       821.0 s          ┘
21:00         1.5 s          ← 阶段3:贴零,管道排空,目标表实时跟随源表

约 5 分钟内从 ~50 分钟积压降到 1.5 秒。 全程 Invocations 11,236 次、Errors 0

讲究

  • 起跳高是正常的(积压时长),不是故障。
  • 多分片必须看 Maximum(最慢分片才是瓶颈),用 Average 会被快分片拉低产生"已追平"假象。
  • IteratorAge 贴零只是"有能力实时追平"的入场券,不是 cutover 终点——源表此刻仍在写。

8.4 关于 Lambda 15 分钟限制与位点(澄清)

  • 15 分钟限制单批:事件源映射由 AWS 托管轮询器循环调用函数,每次只处理 ≤BatchSize 条,几百毫秒返回,每次都是全新 15min 配额。积压再多也是循环调用数万次,碰不到上限。
  • offset/checkpoint 由事件源映射服务端持久化维护,跨函数重启/更新/disable-enable 都不丢,无需自管。位点仅在函数成功返回后推进 → "至少一次"语义 → 必须幂等。
  • 真正的硬天花板是 Stream 24h 保留:追平 < 24h 即安全;大表导入可能超 24h 时改用 Kinesis(保留期最长 1 年)。

9. 阶段八:Cutover 与校验

正确顺序:先停写 → 再确认排空 → 最后校验切换。只看 IteratorAge 低就切会丢数据。

9.1 第 2 步:停止源表写入(停机窗口开始)

停掉模拟客户端进程。

模拟客户端 20:03:42 启动(Export 之前),停写即模拟"上游业务停写或切只读",cutover 停机窗口从此刻开始。

9.2 第 3 步:等管道排空

停写后 IteratorAge 稳定底部(21:00→1.5s, 21:01→1.39s, 21:02→1.42s),Invocations 从峰值 2600+/分 降至 21:02 的 1140,之后归零(21:03+ 无数据点)→ 管道排空LastResult: OK

9.3 第 4 步:数据校验(停写后做才有意义)

(a) 精确条数比对——parallel scan(Select=COUNT,8 段并发,各段循环翻页累加):

def seg(i):
    cnt=0; lek=None
    while True:
        kw=dict(TableName=TABLE, Select='COUNT', Segment=i, TotalSegments=8)
        if lek: kw['ExclusiveStartKey']=lek
        r=ddb.scan(**kw); cnt+=r['Count']; lek=r.get('LastEvaluatedKey')
        if not lek: break
    return cnt
# 源账号与目标账号各跑一次
精确 count
源表 source_account_id1,562,087
目标表 target_account_id1,562,087

精确相等,零误差。 数据守恒:基线 1,504,768 + 增量净增 57,319(INSERT−DELETE)= 1,562,087。

不用 DescribeTable.ItemCount——它 ~6h 才刷新,cutover 当下不准。

(b) 内容级抽样比对——防"删一条又错加一条凑巧 count 相等":

  1. 源账号抽 50 条 metadata.last_op=UPDATE 的记录,存主键 + numeric_field/category/op_ts 到临时文件。
  2. 目标账号逐条 get_item 比对这几个字段。

结果:MATCH 50 / MISMATCH 0 / MISSING 0 —— 全部一致,增量 UPDATE 已正确同步(含改过的字段值)。

9.4 第 5 步及之后

测试环境演练到校验为止即证明方案可行。真实迁移后续:切流量到目标表 → 恢复写入(停机窗口结束)→ 保留源表+Stream 兜底 → 稳定后清理。


10. 全流程踩坑清单(速查)

#现象解法
1Import 压缩类型选 NONEItemValidationError,处理数=文件数,导入 0压缩类型必须选 GZIP
2Stream 不支持 AT_TIMESTAMP无法按时间点起始消费TRIM_HORIZON 全重放 + 代码 EXPORT_TS 过滤
3基线项无 _replication_tsTRIM_HORIZON 重放早于基线的旧记录可能覆盖基线EXPORT_TS 过滤掉早于基线的记录
4IteratorAge 用 Average被快分片拉低,误判已追平看 Maximum
5只看 IteratorAge 低就 cutover源表仍在写,丢增量先停写→再等排空→后校验
6DescribeTable.ItemCount~6h 才刷新,不准parallel scan COUNT 精确统计
7Python print 块缓冲后台输出文件长期 0 字节正常现象;PYTHONUNBUFFERED=1 或读进程状态
8Windows GBK 控制台中文输出乱码PYTHONIOENCODING=utf-8

11. 本次演练创建/变更的资源清单(便于清理)

源账号 source_account_id(cn-north-1)

资源标识操作
DynamoDB 表allinone_qa_LargeDataTable复用,补写 40w 数据
PITR该表开启
Streams该表,NEW_AND_OLD_IMAGES,stream/2026-06-22T12:04:28.095开启
IAM 角色ddb-xacct-replicator-role新建(含内联策略 replicator-inline
Lambda 函数ddb-xacct-replicator新建
事件源映射5706f37f-bcb7-43c0-85c9-6696cb5fbee7新建,当前 Enabled
Export 任务01782130127268-c860c940导出到 menglonq 桶

目标账号 target_account_id(cn-north-1)

资源标识操作
DynamoDB 表allinone_qa_LargeDataTableImport 新建(1,562,087 条)
资源策略RevisionId 1782132298690挂在上表
Import 任务01782130969812-476a6576(成功)、及第一次失败的任务
S3 桶menglonq,prefix tmp/ddb/backup/AWSDynamoDB/01782130127268-c860c940/存导出数据

清理顺序(演练收尾时)

  1. 源账号:disable + 删除事件源映射 5706f37f-...
  2. 源账号:删除 Lambda ddb-xacct-replicator、角色 ddb-xacct-replicator-role
  3. 目标账号:delete_resource_policy 移除目标表资源策略
  4. 源账号:关闭 Streams、按需关闭 PITR
  5. 目标账号:删除导入的测试表(如不再需要)
  6. S3:清理 menglonq/tmp/ddb/backup/ 导出数据
  7. 删除/归档本地脚本产生的临时文件

12. 待改进 / 生产化建议

本次为测试验证,以下在交付客户的生产方案中需补强:

  1. 毒丸防护(必补):当前 MaximumRetryAttempts=-1(无限)且未配 DLQ。生产应在事件源映射配 DestinationConfig.OnFailure 指向 SQS/SNS 死信队列,并设有限 MaximumRetryAttempts(如 3),否则一条永远写不进的毒丸会卡死整个分片直到 24h 过期。
  2. 删除墓碑:原生 Streams 严格保序无重复,正常无乱序删除问题;若改用 Kinesis 承载增量,需额外的删除墓碑机制防"删除后乱序到达旧写入被重新插入"。
  3. 抽样校验工具化:当前抽样比对是临时脚本,生产应做成覆盖更广(含 INSERT/DELETE)、可重复、出报告的校验工具。
  4. cutover 编排:停写、排空确认、校验、切流量、回滚应脚本化编排,减少人工操作窗口。
  5. 回滚机制:cutover 后对目标表也开 Streams 反向捕获或保留短期双写,以便安全回滚。
  6. KMS:若 S3 对象/目标表用客户托管 KMS 密钥,需相应授予解密/加密权限,且密钥与目标 S3 桶同区域;不支持 SSE-C。
  7. 大表选型:导入预计 >24h 时,增量通道改用 Kinesis Data Streams for DynamoDB(保留期最长 1 年)。
  8. 并发提速:若追平慢(IteratorAge 持续涨),调高事件源映射 ParallelizationFactor(默认 1,最大 10)。因用时间戳幂等条件写,并发打破分片内保序也安全。

附:关键时间线

时刻(北京)事件
20:03:42模拟客户端启动(Export 之前),持续写入
20:04:28源表 Streams 开启(先于 Export)
20:06:47EXPORT_TS(exportTime 提前 120s 的过滤起点)
20:08:47Export 完成,exportTime,itemCount 1,504,768
~20:15第一次 Import 失败(压缩 NONE)
~20:22第二次 Import(GZIP)成功,1,504,768 条
~20:54部署 Lambda 三件套(Enabled=False)
~20:54目标表挂资源策略(RevisionId 1782132298690)
20:55enable 事件源映射,IteratorAge ~3037s 起跳
21:00IteratorAge 贴零(1.5s),追平
21:01停模拟客户端(停机窗口开始)
~21:03管道排空
~21:05校验:count 1,562,087==1,562,087,内容 50/50 一致