import argparse import boto3 import json import time from pathlib import Path from botocore.exceptions import ClientError def get_queue_url(sqs_client, queue_name: str) -> str: response = sqs_client.get_queue_url(QueueName=queue_name) return response['QueueUrl'] def sanitize_filename(value: str) -> str: return ''.join(c if c.isalnum() or c in ('-', '_', '.') else '_' for c in value) def save_message(queue_name: str, msg: dict) -> Path: output_dir = Path(queue_name) output_dir.mkdir(parents=True, exist_ok=True) timestamp = msg.get('Attributes', {}).get('SentTimestamp') if timestamp: # SentTimestamp is epoch in milliseconds timestamp = str(timestamp) else: timestamp = str(int(time.time() * 1000)) message_id = sanitize_filename(msg['MessageId']) filename = f'{timestamp}_{message_id}.json' file_path = output_dir / filename payload = { 'MessageId': msg['MessageId'], 'Body': msg['Body'], 'Attributes': msg.get('Attributes', {}), 'MessageAttributes': msg.get('MessageAttributes', {}), 'MD5OfBody': msg.get('MD5OfBody'), } with open(file_path, 'w', encoding='utf-8') as f: json.dump(payload, f, indent=2, ensure_ascii=False) return file_path def poll_existing_messages( sqs_client, queue_name: str, queue_url: str, wait_time_seconds: int = 2, visibility_timeout: int = 30, sleep_seconds: float = 0.2, ) -> int: """Poll all currently visible messages from an SQS queue without deleting them. Important: - Messages become temporarily invisible after being received. - Since we do NOT delete them, they will reappear after visibility timeout. """ total_saved = 0 while True: response = sqs_client.receive_message( QueueUrl=queue_url, MaxNumberOfMessages=10, # SQS max batch size WaitTimeSeconds=wait_time_seconds, VisibilityTimeout=visibility_timeout, AttributeNames=['All'], MessageAttributeNames=['All'], ) messages = response.get('Messages', []) if not messages: break for msg in messages: file_path = save_message(queue_name, msg) total_saved += 1 print(f'Saved: {file_path}') print(f'Downloaded {total_saved} messages so far...') time.sleep(sleep_seconds) return total_saved def main(): parser = argparse.ArgumentParser( description='Poll all currently visible messages from an SQS ' 'queue without deleting them.' ) parser.add_argument( '--sqs-name', default='prod-amazon-data-availability-queue', help='Name of the SQS queue', ) parser.add_argument( '--profile', default='default', help='AWS profile name (optional)', ) parser.add_argument( '--region', default='us-east-1', help='AWS region name (optional)', ) parser.add_argument( '--visibility-timeout', type=int, default=30, help='Visibility timeout in seconds (default: 30)', ) parser.add_argument( '--wait-time-seconds', type=int, default=2, help='Receive wait time in seconds (default: 2)', ) args = parser.parse_args() session_kwargs = {} if args.profile: session_kwargs['profile_name'] = args.profile if args.region: session_kwargs['region_name'] = args.region session = boto3.Session(**session_kwargs) sqs = session.client('sqs') try: queue_url = get_queue_url(sqs, args.sqs_name) print(f'Resolved queue URL: {queue_url}') total_saved = poll_existing_messages( sqs_client=sqs, queue_name=args.sqs_name, queue_url=queue_url, wait_time_seconds=args.wait_time_seconds, visibility_timeout=args.visibility_timeout, ) print(f'Saved {total_saved} messages into folder: {args.sqs_name}') except ClientError as e: print(f'AWS error: {e}') raise SystemExit(1) except Exception as e: print(f'Unexpected error: {e}') raise SystemExit(1) if __name__ == '__main__': main()