You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

百万级查询订阅场景下,新增文档匹配高效通知方案咨询(AWS)

Efficient Subscription Notification System for Book Matches

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

AttributeTypeDescription
subscription_idString (PK)Unique UUID for the subscription (groups denormalized entries)
user_emailStringRecipient email address
frequencyStringimmediate/daily/weekly
categoryString (GSI HASH)Single category from the user's subscription (denormalized)
price_minNumberMinimum price threshold (default: 0)
price_maxNumberMaximum price threshold (default: 99999 for "no upper limit")
name_containsStringOptional keyword for simple name matching (e.g., "Harry Potter")
last_notifiedNumberTimestamp 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:

  1. Queries the ScheduledNotifications table for unprocessed entries
  2. Sends a single summary email per user with all matched books
  3. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 07:38:16