百万级查询订阅场景下,新增文档匹配高效通知方案咨询(AWS)
Based on your requirements (1M+ subscriptions, 100 daily new books, no full-text search needed), here's a scalable, low-overhead solution that avoids iterating through all subscriptions:
Core Idea: Reverse Matching
Instead of checking every subscription against new books, we extract attributes from new books and query only relevant subscriptions using indexed data structures. This cuts down processing from 1M+ checks to a handful of targeted queries per book.
Step 1: Structure Subscriptions for Fast Queries
Store subscriptions in a DynamoDB table with denormalized entries for multi-category subscriptions (each category gets its own row linked to the same subscription ID). This lets us query subscriptions by category directly.
Subscription Table Schema
| Attribute | Type | Description |
|---|---|---|
subscription_id | String (PK) | Unique UUID for the subscription (groups denormalized entries) |
user_email | String | Recipient email address |
frequency | String | immediate/daily/weekly |
category | String (GSI HASH) | Single category from the user's subscription (denormalized) |
price_min | Number | Minimum price threshold (default: 0) |
price_max | Number | Maximum price threshold (default: 99999 for "no upper limit") |
name_contains | String | Optional keyword for simple name matching (e.g., "Harry Potter") |
last_notified | Number | Timestamp of last notification (for deduplication) |
Global Secondary Index (GSI)
Create a GSI named CategoryPriceIndex with:
- Partition key:
category - Sort key:
price_min - Projection:
ALL(includes all attributes)
Step 2: New Book Processing Workflow
Trigger a Lambda from your DynamoDB Stream (when new books are added) to handle matching and notification routing:
Key Logic Snippet (Node.js)
const AWS = require('aws-sdk'); const dynamodb = new AWS.DynamoDB.DocumentClient(); const sqs = new AWS.SQS({ region: 'us-east-1' }); exports.handler = async (event) => { const newBook = event.Records[0].dynamodb.NewImage; const bookCategories = newBook.category.L.map(item => item.S); const bookPrice = parseFloat(newBook.price.N); const bookName = newBook.name.S; // Track unique subscriptions to avoid duplicates const matchedSubs = new Map(); // Query subscriptions for each category in the new book for (const category of bookCategories) { const queryParams = { TableName: 'Subscriptions', IndexName: 'CategoryPriceIndex', KeyConditionExpression: '#cat = :category AND #min <= :price', FilterExpression: '#max >= :price', ExpressionAttributeNames: { '#cat': 'category', '#min': 'price_min', '#max': 'price_max' }, ExpressionAttributeValues: { ':category': category, ':price': bookPrice } }; const results = await dynamodb.query(queryParams).promise(); for (const sub of results.Items) { // Filter name keyword if present if (sub.name_contains && !bookName.includes(sub.name_contains)) continue; // Avoid duplicate subscriptions from multi-category matches if (!matchedSubs.has(sub.subscription_id)) { matchedSubs.set(sub.subscription_id, sub); } } } // Route subscriptions to appropriate notification channels await routeNotifications(matchedSubs, { name: bookName, price: bookPrice, categories: bookCategories }); return { statusCode: 200, body: 'Match processing complete' }; }; async function routeNotifications(subs, book) { const immediateQueueUrl = 'YOUR_SQS_QUEUE_URL'; const dailyBatch = new Map(); const weeklyBatch = new Map(); subs.forEach(sub => { switch (sub.frequency) { case 'immediate': // Send to SQS for async email delivery sqs.sendMessage({ QueueUrl: immediateQueueUrl, MessageBody: JSON.stringify({ user: sub.user_email, book }) }).promise(); break; case 'daily': const dailyKey = `${sub.user_email}#daily#${new Date().toISOString().split('T')[0]}`; dailyBatch.set(dailyKey, [...(dailyBatch.get(dailyKey) || []), book]); break; case 'weekly': const weekStart = new Date(new Date().setDate(new Date().getDate() - new Date().getDay() + 1)); const weeklyKey = `${sub.user_email}#weekly#${weekStart.toISOString().split('T')[0]}`; weeklyBatch.set(weeklyKey, [...(weeklyBatch.get(weeklyKey) || []), book]); break; } }); // Save daily/weekly batches to DynamoDB for scheduled delivery await saveScheduledBatches(dailyBatch, 'daily'); await saveScheduledBatches(weeklyBatch, 'weekly'); } async function saveScheduledBatches(batch, frequency) { if (batch.size === 0) return; const writeRequests = Array.from(batch.entries()).map(([id, books]) => ({ PutRequest: { Item: { id, user_email: id.split('#')[0], frequency, books, processed: false } } })); await dynamodb.batchWrite({ RequestItems: { ScheduledNotifications: writeRequests } }).promise(); }
Step 3: Notification Delivery
Immediate Notifications
Use a separate Lambda to consume the SQS queue and send emails via AWS SES:
const ses = new AWS.SES({ region: 'us-east-1' }); exports.handler = async (event) => { for (const record of event.Records) { const { user, book } = JSON.parse(record.body); await ses.sendEmail({ Destination: { ToAddresses: [user] }, Message: { Subject: { Data: `New Book Match: ${book.name}` }, Body: { Text: { Data: `New book alert:\nName: ${book.name}\nPrice: $${book.price}\nCategories: ${book.categories.join(', ')}` } } }, Source: 'notifications@yourdomain.com' }).promise(); } };
Daily/Weekly Notifications
Use AWS EventBridge to trigger a scheduled Lambda (e.g., daily at 8 AM, weekly on Monday) that:
- Queries the
ScheduledNotificationstable for unprocessed entries - Sends a single summary email per user with all matched books
- Marks entries as processed (or uses DynamoDB TTL to auto-delete old entries)
Why This Works
- No full subscription scans: We only query subscriptions linked to the new book's categories, drastically reducing processing load.
- Scalable: DynamoDB and Lambda handle 1M+ subscriptions and 100 daily books with ease.
- Cost-effective: Minimal query volume and batch processing keep AWS costs low.
- Deduplication: Uses subscription IDs and date-based keys to avoid duplicate notifications.
内容的提问来源于stack exchange,提问作者Supun Induwara

