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

无延迟实时数据处理:基于AWS与Amazon Aurora MySQL的聚合方案咨询

实时数据聚合至Aurora MySQL的最优AWS方案

针对你这个基于AWS Aurora MySQL的实时数据聚合需求——将第一个表的流入数据,实时聚合clicks、impressions等字段到第二个表以支持低延迟查询,我整理了几个适配不同场景的最优技术方案:

方案1:Aurora MySQL 原生触发器(小流量/简单逻辑首选)

如果你的数据流入量不大(比如每秒几十到几百条),且聚合逻辑仅为简单的分组求和,直接用Aurora的数据库触发器是最轻量化的方案。

实现步骤:

  • 在第一个表上创建AFTER INSERT触发器,每次有新数据写入时,自动触发聚合逻辑。
  • 触发器中使用INSERT ... ON DUPLICATE KEY UPDATE语法,对第二个表的聚合记录进行更新(如果维度记录已存在则累加,不存在则插入)。

示例代码:

假设第一个表是raw_events,字段为user_id, ad_id, clicks, impressions, event_time;第二个聚合表是aggregated_stats,主键为user_id, ad_id, date:

DELIMITER //
CREATE TRIGGER update_aggregated_stats AFTER INSERT ON raw_events
FOR EACH ROW
BEGIN
    INSERT INTO aggregated_stats (user_id, ad_id, date, total_clicks, total_impressions)
    VALUES (NEW.user_id, NEW.ad_id, DATE(NEW.event_time), NEW.clicks, NEW.impressions)
    ON DUPLICATE KEY UPDATE
        total_clicks = total_clicks + NEW.clicks,
        total_impressions = total_impressions + NEW.impressions;
END //
DELIMITER ;

优缺点:

  • ✅ 零额外服务成本,实现简单快速
  • ❌ 高流量下会拖慢第一个表的写入性能(触发器同步执行)
  • ❌ 复杂聚合逻辑(比如多窗口、多维度关联)难以实现

方案2:Aurora CDC + AWS Lambda(中等流量/灵活逻辑首选)

当数据流量增大,或者需要更灵活的聚合逻辑时,用Aurora 变更数据捕获(CDC)+ Lambda是更优的选择,它能解耦数据写入和聚合处理,避免影响主表性能。

实现步骤:

  1. 开启Aurora MySQL的CDC功能:可以通过AWS DMS(数据库迁移服务)捕获raw_events表的变更,或者使用Aurora原生的Lambda集成(直接将表的INSERT事件推送到Lambda)。
  2. 在Lambda中编写聚合逻辑:接收CDC事件,按指定维度(比如user_id+ad_id+date)分组,计算clicks和impressions的累加值。
  3. 使用Lambda的批量处理功能:将多个事件合并处理,减少对聚合表的写入次数,提升性能。
  4. 写入聚合表:同样使用INSERT ... ON DUPLICATE KEY UPDATE语法更新aggregated_stats。

示例Lambda代码(Python):

import mysql.connector
from mysql.connector import Error

def lambda_handler(event, context):
    # 从CDC事件中提取批量数据
    records = []
    for record in event['Records']:
        if record['eventName'] == 'INSERT':
            new_image = record['dynamodb']['NewImage']
            records.append({
                'user_id': new_image['user_id']['S'],
                'ad_id': new_image['ad_id']['S'],
                'event_time': new_image['event_time']['S'],
                'clicks': int(new_image['clicks']['N']),
                'impressions': int(new_image['impressions']['N'])
            })
    
    # 按维度分组聚合
    agg_data = {}
    for rec in records:
        date = rec['event_time'].split(' ')[0]
        key = (rec['user_id'], rec['ad_id'], date)
        if key not in agg_data:
            agg_data[key] = {'clicks': 0, 'impressions': 0}
        agg_data[key]['clicks'] += rec['clicks']
        agg_data[key]['impressions'] += rec['impressions']
    
    # 连接Aurora MySQL并写入聚合数据
    try:
        connection = mysql.connector.connect(
            host='your-aurora-endpoint',
            database='your-db',
            user='your-user',
            password='your-password'
        )
        cursor = connection.cursor()
        for (user_id, ad_id, date), stats in agg_data.items():
            query = """
                INSERT INTO aggregated_stats (user_id, ad_id, date, total_clicks, total_impressions)
                VALUES (%s, %s, %s, %s, %s)
                ON DUPLICATE KEY UPDATE
                    total_clicks = total_clicks + %s,
                    total_impressions = total_impressions + %s
            """
            cursor.execute(query, (user_id, ad_id, date, stats['clicks'], stats['impressions'], stats['clicks'], stats['impressions']))
        connection.commit()
    except Error as e:
        print(f"Database error: {e}")
        raise
    finally:
        if connection.is_connected():
            cursor.close()
            connection.close()

优缺点:

  • ✅ 解耦写入和聚合,不影响主表性能
  • ✅ 支持复杂的自定义逻辑(比如多维度关联、时间窗口计算)
  • ✅ 按需付费,成本可控
  • ❌ 需要编写和维护Lambda代码
  • ❌ 超高流量下需注意Lambda的并发限制,可结合SQS做缓冲

方案3:Kinesis Data Streams + Kinesis Data Analytics(高流量/低代码聚合首选)

如果你的数据流入量非常大(每秒数千到数万条),且聚合逻辑可以用SQL表达,那么Kinesis流 + Kinesis Data Analytics是最佳选择,它能实现毫秒级的实时聚合,且无需大量代码开发。

实现步骤:

  1. 调整数据写入路径:移动设备API不再直接写入Aurora,而是将数据发送到Kinesis Data Streams(如果必须先写Aurora,也可以通过Aurora CDC将数据同步到Kinesis流)。
  2. 创建Kinesis Data Analytics应用:使用SQL语句对Kinesis流中的实时数据做聚合,比如按user_id, ad_id和1小时窗口求和。
  3. 输出聚合结果:将Kinesis Data Analytics的聚合结果直接写入Aurora的aggregated_stats表。

示例Kinesis Data Analytics SQL:

CREATE OR REPLACE STREAM AGGREGATED_STREAM (
    user_id VARCHAR(255),
    ad_id VARCHAR(255),
    window_start TIMESTAMP,
    total_clicks INTEGER,
    total_impressions INTEGER
);

CREATE OR REPLACE PUMP AGGREGATION_PUMP AS INSERT INTO AGGREGATED_STREAM
SELECT
    user_id,
    ad_id,
    FLOOR(event_time TO HOUR) AS window_start,
    SUM(clicks) AS total_clicks,
    SUM(impressions) AS total_impressions
FROM
    SOURCE_SQL_STREAM_001
GROUP BY
    user_id, ad_id, FLOOR(event_time TO HOUR);

优缺点:

  • ✅ 支持超高吞吐量,毫秒级延迟
  • ✅ 用SQL实现聚合,低代码/无代码开发
  • ✅ 自动扩容,无需担心并发瓶颈
  • ❌ 相比Lambda方案,自定义逻辑灵活性稍弱
  • ❌ 需要调整数据写入路径(或额外配置CDC到Kinesis)

方案4:AWS Glue Streaming ETL(复杂场景/多数据源聚合首选)

如果你的聚合逻辑需要和其他数据源(比如S3中的历史数据、其他数据库)结合,或者需要更复杂的转换处理,AWS Glue Streaming ETL是合适的选择,它提供了托管的流处理环境,支持Spark Streaming。

实现步骤:

  1. 配置Glue流作业:读取Aurora CDC流或Kinesis流中的数据。
  2. 使用Spark代码编写聚合逻辑:实现多维度关联、复杂窗口计算等。
  3. 将聚合结果写入Aurora的aggregated_stats表。

通用最佳实践:

  • 无论选择哪个方案,aggregated_stats表都应该按聚合维度(比如date)做分区,并设置user_id, ad_id, date为复合主键,确保INSERT ... ON DUPLICATE KEY UPDATE能高效执行。
  • 对于高流量场景,建议开启Aurora的读副本,让聚合查询流量走读副本,进一步降低主库压力。

内容的提问来源于stack exchange,提问作者Osama Alvi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.01 01:52:30