feat(redis): store runs by group then time

This commit is contained in:
tiennm99 committed 2026-06-10 09:47:23 +07:00
1 parent 24bf801656
commit c24b1a9770
3 files changed
+139 -8

No files matched your search

+12 -1
View File
@@ -68,7 +68,7 @@ All data lives in Redis under the `telegram-export` prefix. No key references an
telegram-export:config -> JSON config:
{ api_id, api_hash, phone, group_ids }
telegram-export:session -> StringSession string (login)
telegram-export:run:<yyyymmddhhmmss>:<group_id> -> one group's export as JSON:
telegram-export:run:<group_id>:<yyyymmddhhmmss> -> one group's export as JSON:
{ group_id, title, time,
members: [{ id, username, first_name, last_name }] }
```
@@ -82,6 +82,17 @@ for rec in list_exports(): # sorted by (time, group_id)
print(rec['time'], rec['group_id'], rec['title'], len(rec['members']), 'members')
```
If you have data written with the older `run:<time>:<group_id>` key order,
migrate it once:
```bash
python migrate_run_keys_to_group_time.py --dry-run
python migrate_run_keys_to_group_time.py
```
After verifying compare/history output, run it again with `--delete-old` to
remove the older `run:<time>:<group_id>` keys. Existing new-format keys are kept.
## Rate limits and visibility notes
- Telegram limits how fast you can fetch participants; large groups may take longer.
+21 -7
View File
@@ -2,7 +2,7 @@
Each group's members are stored per export as one self-contained JSON key:
<prefix>:run:<yyyymmddhhmmss>:<group_id>
<prefix>:run:<group_id>:<yyyymmddhhmmss>
-> {group_id, title, time, members:[{id,username,first_name,last_name}]}
All groups in a single run share one timestamp. No key references another, so
@@ -39,7 +39,7 @@ def save_group_export(group_id, title, members, run_time):
'members': [member_dict(m) for m in members],
}
redis_client.set(
key('run', run_time, str(group_id)),
key('run', str(group_id), run_time),
json.dumps(record, ensure_ascii=False),
)
@@ -50,7 +50,9 @@ def list_exports():
Tolerates keys deleted mid-scan and corrupt/non-JSON values (skips them).
"""
records = []
for export_key in redis_client.scan_iter(match=key('run', '*')):
for export_key in redis_client.scan_iter(match=key('run', '*', '*')):
if not _is_group_time_run_key(export_key):
continue
raw = redis_client.get(export_key)
if raw is None: # deleted between scan and get
continue
@@ -73,10 +75,13 @@ def list_group_exports(group_id):
def get_group_export(group_id, run_time):
"""Return a group export at one run time, or None when missing."""
for record in list_group_exports(group_id):
if record.get('time') == run_time:
return record
return None
raw = redis_client.get(key('run', str(group_id), run_time))
if raw is None:
return None
try:
return json.loads(raw)
except (ValueError, TypeError):
return None
def latest_two_group_exports(group_id):
@@ -107,3 +112,12 @@ def _member_index(record):
if member_id is not None:
members[int(member_id)] = member
return members
def _is_run_time(value):
return len(value) == 14 and value.isdigit()
def _is_group_time_run_key(export_key):
parts = export_key.split(':')
return len(parts) == 4 and not _is_run_time(parts[2]) and _is_run_time(parts[3])
+106
View File
@@ -0,0 +1,106 @@
"""Migrate Redis run keys from run:<time>:<group_id> to run:<group_id>:<time>."""
import argparse
import os
import redis
from dotenv import load_dotenv
PREFIX = 'telegram-export'
TIME_LENGTH = 14
def redis_key(*parts):
return ':'.join([PREFIX, *parts])
def parse_args():
parser = argparse.ArgumentParser(
description='Migrate Redis export run keys to group-first format.',
)
parser.add_argument(
'--dry-run',
action='store_true',
help='Print planned changes without writing Redis.',
)
parser.add_argument(
'--overwrite',
action='store_true',
help='Overwrite new-format keys when they already exist.',
)
parser.add_argument(
'--delete-old',
action='store_true',
help='Delete old run:<time>:<group_id> keys after successful copy.',
)
return parser.parse_args()
def connect_redis():
load_dotenv()
redis_url = os.getenv('REDIS_URL')
if not redis_url:
raise SystemExit('REDIS_URL not set in .env')
return redis.from_url(redis_url, decode_responses=True)
def is_run_time(value):
return len(value) == TIME_LENGTH and value.isdigit()
def migrate_run_key(client, old_key, args):
parts = old_key.split(':')
if len(parts) != 4:
return 'skipped'
_, _, run_time, group_id = parts
if not is_run_time(run_time):
return 'skipped'
new_key = redis_key('run', group_id, run_time)
if client.exists(new_key) and not args.overwrite:
if args.delete_old and not args.dry_run:
client.delete(old_key)
print(f'deleted old existing: {old_key}')
return 'deleted'
print(f'skip existing: {new_key}')
return 'skipped'
value = client.get(old_key)
if value is None:
return 'missing'
print(f'{old_key} -> {new_key}')
if args.dry_run:
return 'planned'
client.set(new_key, value)
if args.delete_old:
client.delete(old_key)
return 'migrated'
def main():
args = parse_args()
client = connect_redis()
counts = {'planned': 0, 'migrated': 0, 'deleted': 0, 'skipped': 0, 'missing': 0}
for old_key in client.scan_iter(match=redis_key('run', '*', '*')):
result = migrate_run_key(client, old_key, args)
counts[result] += 1
print('')
print(
'runs: '
f'{counts["migrated"]} migrated, '
f'{counts["deleted"]} deleted, '
f'{counts["planned"]} planned, '
f'{counts["skipped"]} skipped, '
f'{counts["missing"]} missing'
)
if args.dry_run:
print('dry run only; no Redis keys changed.')
if __name__ == '__main__':
raise SystemExit(main())