Skip to content

Commit 0e43d27

Browse files
committed
fix: flush activities and checkpoint heads in batches during onboarding (CM-1318)
Signed-off-by: Uroš Marolt <uros@marolt.me>
1 parent bfdcf1e commit 0e43d27

2 files changed

Lines changed: 14 additions & 2 deletions

File tree

services/apps/mailing_list_integration/src/crowdmail/settings.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,3 +41,8 @@ def load_env_var(key: str, required=True, default=None):
4141
STUCK_RECURRENT_LIST_TIMEOUT_HOURS = int(
4242
load_env_var("STUCK_RECURRENT_LIST_TIMEOUT_HOURS", default="4")
4343
)
44+
45+
# Flush accumulated activities (and checkpoint processed heads) every this many
46+
# messages instead of buffering an entire list's history in memory before one
47+
# flush at the end — large lore archives can have 100k+ messages.
48+
ACTIVITY_FLUSH_BATCH_SIZE = int(load_env_var("ACTIVITY_FLUSH_BATCH_SIZE", default="500"))

services/apps/mailing_list_integration/src/crowdmail/worker/list_worker.py

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
from crowdmail.services.parse.noteren import parse_email
2222
from crowdmail.services.queue import queue_service
2323
from crowdmail.settings import (
24+
ACTIVITY_FLUSH_BATCH_SIZE,
2425
DEFAULT_TENANT_ID,
2526
WORKER_ERROR_BACKOFF_SEC,
2627
WORKER_POLLING_INTERVAL_SEC,
@@ -90,6 +91,7 @@ async def _process_single_list(self, mailing_list: MailingList):
9091
shard = shard_index(shard_path)
9192
commit_ids = await new_commits(shard_path, heads.get(shard))
9293
for git_id in commit_ids:
94+
heads[shard] = git_id
9395
try:
9496
message, blob_id = read_email(shard_path, git_id)
9597
parsed = parse_email(
@@ -134,8 +136,13 @@ async def _process_single_list(self, mailing_list: MailingList):
134136
mailing_list.segment_id, mailing_list.integration_id, result_id
135137
)
136138
)
137-
if commit_ids:
138-
heads[shard] = commit_ids[-1]
139+
140+
if len(activities_db) >= ACTIVITY_FLUSH_BATCH_SIZE:
141+
await batch_insert_activities(activities_db)
142+
await queue_service.send_batch_activities(activities_kafka)
143+
await update_processed_heads(mailing_list.id, heads)
144+
activities_db = []
145+
activities_kafka = []
139146

140147
if activities_db:
141148
await batch_insert_activities(activities_db)

0 commit comments

Comments
 (0)