mirror of
https://github.com/fawney19/Aether.git
synced 2026-09-01 17:00:21 +08:00
refactor(maintenance): body 压缩改为逐条独立事务,降低内存占用与锁粒度
将 _cleanup_body_fields 从批量加载完整记录改为先查询 ID 列表, 再逐条独立会话处理压缩,避免大批量事务导致的内存和锁问题。 批次大小上限从 100 降至 25,新增排序保证处理顺序确定性。
This commit is contained in:
@@ -1025,117 +1025,29 @@ class MaintenanceScheduler:
|
|||||||
|
|
||||||
total_compressed = 0
|
total_compressed = 0
|
||||||
no_progress_count = 0
|
no_progress_count = 0
|
||||||
memory_safe_batch_size = max(1, min(batch_size, 100))
|
memory_safe_batch_size = max(1, min(batch_size, 25))
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
batch_db = create_session()
|
batch_db = create_session()
|
||||||
try:
|
try:
|
||||||
records = (
|
record_ids = [
|
||||||
batch_db.query(
|
row.id
|
||||||
Usage.id,
|
for row in (
|
||||||
Usage.request_body,
|
batch_db.query(Usage.id)
|
||||||
Usage.response_body,
|
.filter(Usage.created_at < cutoff_time)
|
||||||
Usage.provider_request_body,
|
.filter(
|
||||||
Usage.client_response_body,
|
(Usage.request_body.isnot(None))
|
||||||
|
| (Usage.response_body.isnot(None))
|
||||||
|
| (Usage.provider_request_body.isnot(None))
|
||||||
|
| (Usage.client_response_body.isnot(None))
|
||||||
|
)
|
||||||
|
.order_by(Usage.created_at.asc(), Usage.id.asc())
|
||||||
|
.limit(memory_safe_batch_size)
|
||||||
|
.all()
|
||||||
)
|
)
|
||||||
.filter(Usage.created_at < cutoff_time)
|
|
||||||
.filter(
|
|
||||||
(Usage.request_body.isnot(None))
|
|
||||||
| (Usage.response_body.isnot(None))
|
|
||||||
| (Usage.provider_request_body.isnot(None))
|
|
||||||
| (Usage.client_response_body.isnot(None))
|
|
||||||
)
|
|
||||||
.limit(memory_safe_batch_size)
|
|
||||||
.all()
|
|
||||||
)
|
|
||||||
|
|
||||||
if not records:
|
|
||||||
break
|
|
||||||
|
|
||||||
valid_records = [
|
|
||||||
r
|
|
||||||
for r in records
|
|
||||||
if r.request_body is not None
|
|
||||||
or r.response_body is not None
|
|
||||||
or r.provider_request_body is not None
|
|
||||||
or r.client_response_body is not None
|
|
||||||
]
|
]
|
||||||
|
|
||||||
if not valid_records:
|
|
||||||
logger.warning(
|
|
||||||
f"检测到 {len(records)} 条记录的 body 字段为 JSON null,进行清理"
|
|
||||||
)
|
|
||||||
for r in records:
|
|
||||||
batch_db.execute(
|
|
||||||
update(Usage)
|
|
||||||
.where(Usage.id == r.id)
|
|
||||||
.values(
|
|
||||||
request_body=null(),
|
|
||||||
response_body=null(),
|
|
||||||
provider_request_body=null(),
|
|
||||||
client_response_body=null(),
|
|
||||||
)
|
|
||||||
)
|
|
||||||
batch_db.commit()
|
|
||||||
continue
|
|
||||||
|
|
||||||
batch_success = 0
|
|
||||||
|
|
||||||
for r in valid_records:
|
|
||||||
try:
|
|
||||||
result = batch_db.execute(
|
|
||||||
update(Usage)
|
|
||||||
.where(Usage.id == r.id)
|
|
||||||
.values(
|
|
||||||
request_body=null(),
|
|
||||||
response_body=null(),
|
|
||||||
provider_request_body=null(),
|
|
||||||
client_response_body=null(),
|
|
||||||
request_body_compressed=(
|
|
||||||
compress_json(r.request_body) if r.request_body else None
|
|
||||||
),
|
|
||||||
response_body_compressed=(
|
|
||||||
compress_json(r.response_body) if r.response_body else None
|
|
||||||
),
|
|
||||||
provider_request_body_compressed=(
|
|
||||||
compress_json(r.provider_request_body)
|
|
||||||
if r.provider_request_body
|
|
||||||
else None
|
|
||||||
),
|
|
||||||
client_response_body_compressed=(
|
|
||||||
compress_json(r.client_response_body)
|
|
||||||
if r.client_response_body
|
|
||||||
else None
|
|
||||||
),
|
|
||||||
)
|
|
||||||
.execution_options(synchronize_session=False)
|
|
||||||
)
|
|
||||||
if result.rowcount > 0:
|
|
||||||
batch_success += 1
|
|
||||||
except Exception as e:
|
|
||||||
logger.warning(f"压缩记录 {r.id} 失败: {e}")
|
|
||||||
continue
|
|
||||||
|
|
||||||
batch_db.commit()
|
|
||||||
|
|
||||||
if batch_success == 0:
|
|
||||||
no_progress_count += 1
|
|
||||||
if no_progress_count >= 3:
|
|
||||||
logger.error(
|
|
||||||
f"压缩 body 字段连续 {no_progress_count} 批无进展,"
|
|
||||||
"终止循环以避免死循环"
|
|
||||||
)
|
|
||||||
break
|
|
||||||
else:
|
|
||||||
no_progress_count = 0
|
|
||||||
|
|
||||||
total_compressed += batch_success
|
|
||||||
logger.debug(
|
|
||||||
f"已压缩 {batch_success} 条记录的 body 字段,累计 {total_compressed} 条"
|
|
||||||
)
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"压缩 body 字段失败: {e}")
|
logger.exception("加载待压缩 body 记录 ID 失败: {}", e)
|
||||||
try:
|
try:
|
||||||
batch_db.rollback()
|
batch_db.rollback()
|
||||||
except Exception:
|
except Exception:
|
||||||
@@ -1144,6 +1056,114 @@ class MaintenanceScheduler:
|
|||||||
finally:
|
finally:
|
||||||
batch_db.close()
|
batch_db.close()
|
||||||
|
|
||||||
|
if not record_ids:
|
||||||
|
break
|
||||||
|
|
||||||
|
batch_success = 0
|
||||||
|
batch_progress = False
|
||||||
|
|
||||||
|
for record_id in record_ids:
|
||||||
|
record_db = create_session()
|
||||||
|
try:
|
||||||
|
record = (
|
||||||
|
record_db.query(
|
||||||
|
Usage.id,
|
||||||
|
Usage.request_body,
|
||||||
|
Usage.response_body,
|
||||||
|
Usage.provider_request_body,
|
||||||
|
Usage.client_response_body,
|
||||||
|
)
|
||||||
|
.filter(Usage.id == record_id)
|
||||||
|
.first()
|
||||||
|
)
|
||||||
|
|
||||||
|
if record is None:
|
||||||
|
continue
|
||||||
|
|
||||||
|
has_body_payload = (
|
||||||
|
record.request_body is not None
|
||||||
|
or record.response_body is not None
|
||||||
|
or record.provider_request_body is not None
|
||||||
|
or record.client_response_body is not None
|
||||||
|
)
|
||||||
|
|
||||||
|
if not has_body_payload:
|
||||||
|
result = record_db.execute(
|
||||||
|
update(Usage)
|
||||||
|
.where(Usage.id == record.id)
|
||||||
|
.values(
|
||||||
|
request_body=null(),
|
||||||
|
response_body=null(),
|
||||||
|
provider_request_body=null(),
|
||||||
|
client_response_body=null(),
|
||||||
|
)
|
||||||
|
.execution_options(synchronize_session=False)
|
||||||
|
)
|
||||||
|
record_db.commit()
|
||||||
|
if result.rowcount > 0:
|
||||||
|
batch_progress = True
|
||||||
|
continue
|
||||||
|
|
||||||
|
result = record_db.execute(
|
||||||
|
update(Usage)
|
||||||
|
.where(Usage.id == record.id)
|
||||||
|
.values(
|
||||||
|
request_body=null(),
|
||||||
|
response_body=null(),
|
||||||
|
provider_request_body=null(),
|
||||||
|
client_response_body=null(),
|
||||||
|
request_body_compressed=(
|
||||||
|
compress_json(record.request_body) if record.request_body else None
|
||||||
|
),
|
||||||
|
response_body_compressed=(
|
||||||
|
compress_json(record.response_body)
|
||||||
|
if record.response_body
|
||||||
|
else None
|
||||||
|
),
|
||||||
|
provider_request_body_compressed=(
|
||||||
|
compress_json(record.provider_request_body)
|
||||||
|
if record.provider_request_body
|
||||||
|
else None
|
||||||
|
),
|
||||||
|
client_response_body_compressed=(
|
||||||
|
compress_json(record.client_response_body)
|
||||||
|
if record.client_response_body
|
||||||
|
else None
|
||||||
|
),
|
||||||
|
)
|
||||||
|
.execution_options(synchronize_session=False)
|
||||||
|
)
|
||||||
|
record_db.commit()
|
||||||
|
|
||||||
|
if result.rowcount > 0:
|
||||||
|
batch_success += 1
|
||||||
|
batch_progress = True
|
||||||
|
except Exception as e:
|
||||||
|
logger.warning("压缩记录 {} 失败: {}", record_id, e)
|
||||||
|
try:
|
||||||
|
record_db.rollback()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
continue
|
||||||
|
finally:
|
||||||
|
record_db.close()
|
||||||
|
|
||||||
|
if not batch_progress:
|
||||||
|
no_progress_count += 1
|
||||||
|
if no_progress_count >= 3:
|
||||||
|
logger.error(
|
||||||
|
f"压缩 body 字段连续 {no_progress_count} 批无进展," "终止循环以避免死循环"
|
||||||
|
)
|
||||||
|
break
|
||||||
|
else:
|
||||||
|
no_progress_count = 0
|
||||||
|
|
||||||
|
total_compressed += batch_success
|
||||||
|
if batch_success > 0:
|
||||||
|
logger.debug(
|
||||||
|
f"已压缩 {batch_success} 条记录的 body 字段,累计 {total_compressed} 条"
|
||||||
|
)
|
||||||
|
|
||||||
return total_compressed
|
return total_compressed
|
||||||
|
|
||||||
def _cleanup_compressed_fields(self, cutoff_time: datetime, batch_size: int) -> int:
|
def _cleanup_compressed_fields(self, cutoff_time: datetime, batch_size: int) -> int:
|
||||||
|
|||||||
@@ -122,3 +122,102 @@ async def test_candidate_cleanup_uses_dedicated_retention_and_batch_settings(
|
|||||||
assert batch_one.committed is True
|
assert batch_one.committed is True
|
||||||
assert batch_one.closed is True
|
assert batch_one.closed is True
|
||||||
assert batch_two.closed is True
|
assert batch_two.closed is True
|
||||||
|
|
||||||
|
|
||||||
|
def test_cleanup_body_fields_loads_ids_then_processes_records_individually(
|
||||||
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
) -> None:
|
||||||
|
scheduler = MaintenanceScheduler()
|
||||||
|
|
||||||
|
class _IdBatchSession:
|
||||||
|
def __init__(self, ids: list[str]) -> None:
|
||||||
|
self.ids = ids
|
||||||
|
self.closed = False
|
||||||
|
self.query_obj = MagicMock()
|
||||||
|
filtered = self.query_obj.filter.return_value
|
||||||
|
filtered.filter.return_value.order_by.return_value.limit.return_value.all.return_value = [
|
||||||
|
SimpleNamespace(id=value) for value in ids
|
||||||
|
]
|
||||||
|
|
||||||
|
def query(self, *args): # type: ignore[no-untyped-def]
|
||||||
|
self.query_args = args
|
||||||
|
return self.query_obj
|
||||||
|
|
||||||
|
def rollback(self) -> None:
|
||||||
|
raise AssertionError("rollback should not be called")
|
||||||
|
|
||||||
|
def close(self) -> None:
|
||||||
|
self.closed = True
|
||||||
|
|
||||||
|
class _RecordSession:
|
||||||
|
def __init__(self, record: SimpleNamespace) -> None:
|
||||||
|
self.record = record
|
||||||
|
self.closed = False
|
||||||
|
self.committed = False
|
||||||
|
self.query_obj = MagicMock()
|
||||||
|
self.query_obj.filter.return_value.first.return_value = record
|
||||||
|
|
||||||
|
def query(self, *args): # type: ignore[no-untyped-def]
|
||||||
|
self.query_args = args
|
||||||
|
return self.query_obj
|
||||||
|
|
||||||
|
def execute(self, _statement): # type: ignore[no-untyped-def]
|
||||||
|
return SimpleNamespace(rowcount=1)
|
||||||
|
|
||||||
|
def commit(self) -> None:
|
||||||
|
self.committed = True
|
||||||
|
|
||||||
|
def rollback(self) -> None:
|
||||||
|
raise AssertionError("rollback should not be called")
|
||||||
|
|
||||||
|
def close(self) -> None:
|
||||||
|
self.closed = True
|
||||||
|
|
||||||
|
batch_one = _IdBatchSession(["usage-1", "usage-2"])
|
||||||
|
record_one = _RecordSession(
|
||||||
|
SimpleNamespace(
|
||||||
|
id="usage-1",
|
||||||
|
request_body={"hello": "world"},
|
||||||
|
response_body=None,
|
||||||
|
provider_request_body=None,
|
||||||
|
client_response_body=None,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
record_two = _RecordSession(
|
||||||
|
SimpleNamespace(
|
||||||
|
id="usage-2",
|
||||||
|
request_body=None,
|
||||||
|
response_body={"ok": True},
|
||||||
|
provider_request_body=None,
|
||||||
|
client_response_body=None,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
batch_two = _IdBatchSession([])
|
||||||
|
sessions = iter([batch_one, record_one, record_two, batch_two])
|
||||||
|
|
||||||
|
monkeypatch.setattr(
|
||||||
|
maintenance_scheduler_module,
|
||||||
|
"create_session",
|
||||||
|
lambda: next(sessions),
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
maintenance_scheduler_module,
|
||||||
|
"compress_json",
|
||||||
|
lambda payload: f"compressed:{payload}".encode(),
|
||||||
|
)
|
||||||
|
|
||||||
|
compressed = scheduler._cleanup_body_fields(
|
||||||
|
cutoff_time=SimpleNamespace(), # type: ignore[arg-type]
|
||||||
|
batch_size=1000,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert compressed == 2
|
||||||
|
assert batch_one.query_args == (maintenance_scheduler_module.Usage.id,)
|
||||||
|
assert len(record_one.query_args) == 5
|
||||||
|
assert len(record_two.query_args) == 5
|
||||||
|
assert batch_one.closed is True
|
||||||
|
assert batch_two.closed is True
|
||||||
|
assert record_one.committed is True
|
||||||
|
assert record_two.committed is True
|
||||||
|
assert record_one.closed is True
|
||||||
|
assert record_two.closed is True
|
||||||
|
|||||||
Reference in New Issue
Block a user