fix(deploy): migrate 跨 pod 串行化 + rescale 迁移幂等守卫(测试环境×10⁶事故复盘)
事故:连推 3 commit 触发 3 轮 rolling 部署,多个新 pod 并发跑 migrate 且崩溃重跑,
积分 ×10 rescale 被交错重放——ai.0025 执行 6 次(单价/任务计价 ×10⁶),accounts.0008
执行 2 次(限额 ×100),billing.0003/0004 从未完成(django_migrations 漏记录)。
测试库数据已按精确倍率手工修复并补记迁移记录(备份于本机)。
两层防复发:
1. docker-entrypoint 用 MySQL GET_LOCK('airshelf_migrate') 串行化 migrate,
后到 pod 等锁,拿到时迁移已被记录 → 自然 no-op;
2. 三个 rescale 迁移加 airshelf_rescale_marker 幂等标记(与数据变更同事务提交):
记录丢失/崩溃重跑时,标记在 → 跳过,不会重复 ×10。
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -23,10 +23,41 @@ REVERSE_STATEMENTS = [
|
|||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
|
# 幂等守卫(2026-07-07 测试环境事故复盘):并发/崩溃后的 migrate 重跑会把 ×10 重复应用
|
||||||
|
# (事故中本迁移被交错重放,单价被乘到 ×10⁶)。标记行与数据变更同一事务提交:
|
||||||
|
# 上次成功 → 标记在 → 跳过;上次崩溃回滚 → 标记不在 → 安全重放。
|
||||||
|
MARKER = "accounts.0008_points_rescale"
|
||||||
|
_MARKER_DDL = "CREATE TABLE IF NOT EXISTS airshelf_rescale_marker (name varchar(80) NOT NULL PRIMARY KEY)"
|
||||||
|
|
||||||
|
|
||||||
def _run(statements):
|
def _run(statements):
|
||||||
def apply(apps, schema_editor):
|
def apply(apps, schema_editor):
|
||||||
with transaction.atomic(using=schema_editor.connection.alias):
|
conn = schema_editor.connection
|
||||||
with schema_editor.connection.cursor() as cursor:
|
with conn.cursor() as cursor: # DDL 幂等,MySQL 隐式提交故放事务外
|
||||||
|
cursor.execute(_MARKER_DDL)
|
||||||
|
with transaction.atomic(using=conn.alias):
|
||||||
|
with conn.cursor() as cursor:
|
||||||
|
cursor.execute("SELECT COUNT(*) FROM airshelf_rescale_marker WHERE name = %s", [MARKER])
|
||||||
|
if cursor.fetchone()[0]:
|
||||||
|
return # 已应用过(django_migrations 记录丢失/并发重跑),幂等跳过
|
||||||
|
cursor.execute("INSERT INTO airshelf_rescale_marker (name) VALUES (%s)", [MARKER])
|
||||||
|
for sql in statements:
|
||||||
|
cursor.execute(sql)
|
||||||
|
|
||||||
|
return apply
|
||||||
|
|
||||||
|
|
||||||
|
def _run_reverse(statements):
|
||||||
|
def apply(apps, schema_editor):
|
||||||
|
conn = schema_editor.connection
|
||||||
|
with conn.cursor() as cursor:
|
||||||
|
cursor.execute(_MARKER_DDL)
|
||||||
|
with transaction.atomic(using=conn.alias):
|
||||||
|
with conn.cursor() as cursor:
|
||||||
|
cursor.execute("SELECT COUNT(*) FROM airshelf_rescale_marker WHERE name = %s", [MARKER])
|
||||||
|
if not cursor.fetchone()[0]:
|
||||||
|
return # 未应用过或已回滚,无需反向
|
||||||
|
cursor.execute("DELETE FROM airshelf_rescale_marker WHERE name = %s", [MARKER])
|
||||||
for sql in statements:
|
for sql in statements:
|
||||||
cursor.execute(sql)
|
cursor.execute(sql)
|
||||||
|
|
||||||
@@ -35,4 +66,4 @@ def _run(statements):
|
|||||||
|
|
||||||
class Migration(migrations.Migration):
|
class Migration(migrations.Migration):
|
||||||
dependencies = [("accounts", "0007_team_monthly_credit_limit")]
|
dependencies = [("accounts", "0007_team_monthly_credit_limit")]
|
||||||
operations = [migrations.RunPython(_run(RESCALE_STATEMENTS), _run(REVERSE_STATEMENTS))]
|
operations = [migrations.RunPython(_run(RESCALE_STATEMENTS), _run_reverse(REVERSE_STATEMENTS))]
|
||||||
|
|||||||
@@ -20,10 +20,41 @@ REVERSE_STATEMENTS = [
|
|||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
|
# 幂等守卫(2026-07-07 测试环境事故复盘):并发/崩溃后的 migrate 重跑会把 ×10 重复应用
|
||||||
|
# (事故中本迁移被交错重放,单价被乘到 ×10⁶)。标记行与数据变更同一事务提交:
|
||||||
|
# 上次成功 → 标记在 → 跳过;上次崩溃回滚 → 标记不在 → 安全重放。
|
||||||
|
MARKER = "ai.0025_points_rescale"
|
||||||
|
_MARKER_DDL = "CREATE TABLE IF NOT EXISTS airshelf_rescale_marker (name varchar(80) NOT NULL PRIMARY KEY)"
|
||||||
|
|
||||||
|
|
||||||
def _run(statements):
|
def _run(statements):
|
||||||
def apply(apps, schema_editor):
|
def apply(apps, schema_editor):
|
||||||
with transaction.atomic(using=schema_editor.connection.alias):
|
conn = schema_editor.connection
|
||||||
with schema_editor.connection.cursor() as cursor:
|
with conn.cursor() as cursor: # DDL 幂等,MySQL 隐式提交故放事务外
|
||||||
|
cursor.execute(_MARKER_DDL)
|
||||||
|
with transaction.atomic(using=conn.alias):
|
||||||
|
with conn.cursor() as cursor:
|
||||||
|
cursor.execute("SELECT COUNT(*) FROM airshelf_rescale_marker WHERE name = %s", [MARKER])
|
||||||
|
if cursor.fetchone()[0]:
|
||||||
|
return # 已应用过(django_migrations 记录丢失/并发重跑),幂等跳过
|
||||||
|
cursor.execute("INSERT INTO airshelf_rescale_marker (name) VALUES (%s)", [MARKER])
|
||||||
|
for sql in statements:
|
||||||
|
cursor.execute(sql)
|
||||||
|
|
||||||
|
return apply
|
||||||
|
|
||||||
|
|
||||||
|
def _run_reverse(statements):
|
||||||
|
def apply(apps, schema_editor):
|
||||||
|
conn = schema_editor.connection
|
||||||
|
with conn.cursor() as cursor:
|
||||||
|
cursor.execute(_MARKER_DDL)
|
||||||
|
with transaction.atomic(using=conn.alias):
|
||||||
|
with conn.cursor() as cursor:
|
||||||
|
cursor.execute("SELECT COUNT(*) FROM airshelf_rescale_marker WHERE name = %s", [MARKER])
|
||||||
|
if not cursor.fetchone()[0]:
|
||||||
|
return # 未应用过或已回滚,无需反向
|
||||||
|
cursor.execute("DELETE FROM airshelf_rescale_marker WHERE name = %s", [MARKER])
|
||||||
for sql in statements:
|
for sql in statements:
|
||||||
cursor.execute(sql)
|
cursor.execute(sql)
|
||||||
|
|
||||||
@@ -32,4 +63,4 @@ def _run(statements):
|
|||||||
|
|
||||||
class Migration(migrations.Migration):
|
class Migration(migrations.Migration):
|
||||||
dependencies = [("ai", "0024_aitask_base_cost")]
|
dependencies = [("ai", "0024_aitask_base_cost")]
|
||||||
operations = [migrations.RunPython(_run(RESCALE_STATEMENTS), _run(REVERSE_STATEMENTS))]
|
operations = [migrations.RunPython(_run(RESCALE_STATEMENTS), _run_reverse(REVERSE_STATEMENTS))]
|
||||||
|
|||||||
@@ -29,10 +29,41 @@ REVERSE_STATEMENTS = [
|
|||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
|
# 幂等守卫(2026-07-07 测试环境事故复盘):并发/崩溃后的 migrate 重跑会把 ×10 重复应用
|
||||||
|
# (事故中本迁移被交错重放,单价被乘到 ×10⁶)。标记行与数据变更同一事务提交:
|
||||||
|
# 上次成功 → 标记在 → 跳过;上次崩溃回滚 → 标记不在 → 安全重放。
|
||||||
|
MARKER = "billing.0004_points_rescale"
|
||||||
|
_MARKER_DDL = "CREATE TABLE IF NOT EXISTS airshelf_rescale_marker (name varchar(80) NOT NULL PRIMARY KEY)"
|
||||||
|
|
||||||
|
|
||||||
def _run(statements):
|
def _run(statements):
|
||||||
def apply(apps, schema_editor):
|
def apply(apps, schema_editor):
|
||||||
with transaction.atomic(using=schema_editor.connection.alias):
|
conn = schema_editor.connection
|
||||||
with schema_editor.connection.cursor() as cursor:
|
with conn.cursor() as cursor: # DDL 幂等,MySQL 隐式提交故放事务外
|
||||||
|
cursor.execute(_MARKER_DDL)
|
||||||
|
with transaction.atomic(using=conn.alias):
|
||||||
|
with conn.cursor() as cursor:
|
||||||
|
cursor.execute("SELECT COUNT(*) FROM airshelf_rescale_marker WHERE name = %s", [MARKER])
|
||||||
|
if cursor.fetchone()[0]:
|
||||||
|
return # 已应用过(django_migrations 记录丢失/并发重跑),幂等跳过
|
||||||
|
cursor.execute("INSERT INTO airshelf_rescale_marker (name) VALUES (%s)", [MARKER])
|
||||||
|
for sql in statements:
|
||||||
|
cursor.execute(sql)
|
||||||
|
|
||||||
|
return apply
|
||||||
|
|
||||||
|
|
||||||
|
def _run_reverse(statements):
|
||||||
|
def apply(apps, schema_editor):
|
||||||
|
conn = schema_editor.connection
|
||||||
|
with conn.cursor() as cursor:
|
||||||
|
cursor.execute(_MARKER_DDL)
|
||||||
|
with transaction.atomic(using=conn.alias):
|
||||||
|
with conn.cursor() as cursor:
|
||||||
|
cursor.execute("SELECT COUNT(*) FROM airshelf_rescale_marker WHERE name = %s", [MARKER])
|
||||||
|
if not cursor.fetchone()[0]:
|
||||||
|
return # 未应用过或已回滚,无需反向
|
||||||
|
cursor.execute("DELETE FROM airshelf_rescale_marker WHERE name = %s", [MARKER])
|
||||||
for sql in statements:
|
for sql in statements:
|
||||||
cursor.execute(sql)
|
cursor.execute(sql)
|
||||||
|
|
||||||
@@ -41,4 +72,4 @@ def _run(statements):
|
|||||||
|
|
||||||
class Migration(migrations.Migration):
|
class Migration(migrations.Migration):
|
||||||
dependencies = [("billing", "0003_billingconfig")]
|
dependencies = [("billing", "0003_billingconfig")]
|
||||||
operations = [migrations.RunPython(_run(RESCALE_STATEMENTS), _run(REVERSE_STATEMENTS))]
|
operations = [migrations.RunPython(_run(RESCALE_STATEMENTS), _run_reverse(REVERSE_STATEMENTS))]
|
||||||
|
|||||||
@@ -3,10 +3,39 @@ set -e
|
|||||||
|
|
||||||
# Only the web (gunicorn) container should run migrations / collectstatic.
|
# Only the web (gunicorn) container should run migrations / collectstatic.
|
||||||
# The celery worker shares this image but skips DB schema mutation to avoid races.
|
# The celery worker shares this image but skips DB schema mutation to avoid races.
|
||||||
|
#
|
||||||
|
# ⚠️ migrate 必须跨 pod 串行化(MySQL GET_LOCK):2026-07-07 事故——连推 3 个 commit 触发
|
||||||
|
# 3 轮 rolling 部署,多个新 pod 并发跑 migrate,数据迁移(积分 ×10 rescale)被交错重放 6 次,
|
||||||
|
# 测试库单价被乘成 ×10⁶(¥2 → 2,000,000)。锁把并发压成串行;后到者拿到锁时迁移已被
|
||||||
|
# 先到者记录,migrate 自然 no-op。锁超时 600s 拿不到 → 快速失败重启,绝不裸跑。
|
||||||
case "$1" in
|
case "$1" in
|
||||||
gunicorn)
|
gunicorn)
|
||||||
echo "[entrypoint] running migrations..."
|
echo "[entrypoint] running migrations (serialized via DB advisory lock)..."
|
||||||
python manage.py migrate --noinput
|
python - <<'PYEOF'
|
||||||
|
import os
|
||||||
|
|
||||||
|
import django
|
||||||
|
|
||||||
|
os.environ.setdefault("DJANGO_SETTINGS_MODULE", "airshelf.settings.production")
|
||||||
|
django.setup()
|
||||||
|
|
||||||
|
from django.core.management import call_command
|
||||||
|
from django.db import connection
|
||||||
|
|
||||||
|
if connection.vendor == "mysql":
|
||||||
|
with connection.cursor() as cursor:
|
||||||
|
cursor.execute("SELECT GET_LOCK('airshelf_migrate', 600)")
|
||||||
|
acquired = cursor.fetchone()[0]
|
||||||
|
if not acquired:
|
||||||
|
raise SystemExit("[entrypoint] FATAL: migrate advisory lock timeout (another pod stuck?)")
|
||||||
|
try:
|
||||||
|
call_command("migrate", interactive=False)
|
||||||
|
finally:
|
||||||
|
with connection.cursor() as cursor:
|
||||||
|
cursor.execute("SELECT RELEASE_LOCK('airshelf_migrate')")
|
||||||
|
else:
|
||||||
|
call_command("migrate", interactive=False)
|
||||||
|
PYEOF
|
||||||
echo "[entrypoint] collecting static..."
|
echo "[entrypoint] collecting static..."
|
||||||
python manage.py collectstatic --noinput
|
python manage.py collectstatic --noinput
|
||||||
;;
|
;;
|
||||||
|
|||||||
Reference in New Issue
Block a user