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

如何统计CloudFront HLS流媒体每用户数据传输量并存储至MongoDB

Hey there! Let's break down how to track per-user data consumption (in MB) for your HLS streaming workflow and store those stats in MongoDB. Here's a practical, step-by-step approach tailored to your architecture:

Step 1: Configure CloudFront to Capture Critical Metadata

First, you need CloudFront to log the data points that let you tie streaming requests to specific users and videos. Here's what to set up:

  • Enable CloudFront Access Logs: Route logs to a dedicated S3 bucket (separate from your media bucket). These logs are stored as gzipped TSV files.
  • Include User & Video Identifiers:
    • Have your mobile app send a custom HTTP header like X-User-ID with every streaming request (make sure this ID is validated server-side to prevent tampering).
    • In your CloudFront distribution's behavior settings, add X-User-ID to the "Include Headers" list so it shows up in logs.
    • Ensure your HLS file paths include a unique video ID (e.g., /transcoded/{video-id}/master.m3u8)—this lets you link requests to specific videos.

CloudFront logs will now include:

  • sc-bytes: The total bytes sent to the client for each request
  • cs-uri-stem: The path of the requested HLS file (m3u8/ts)
  • x-edge-request-header-X-User-ID: The user ID from the custom header
Step 2: Automate Log Processing with Lambda

Set up a Lambda function triggered by new log files being uploaded to your CloudFront log bucket. This function will parse logs, calculate MB usage, and update MongoDB.

Here's a simplified Python example of the Lambda logic:

import boto3
import gzip
import csv
from pymongo import MongoClient
from datetime import datetime
import os

def lambda_handler(event, context):
    # Initialize clients
    s3 = boto3.client('s3')
    mongo_uri = os.environ['MONGO_URI']
    client = MongoClient(mongo_uri)
    db = client['streaming_stats']
    consumption_coll = db['user_video_consumption']

    # Fetch and decompress the log file
    record = event['Records'][0]['s3']
    bucket = record['bucket']['name']
    key = record['object']['key']
    
    try:
        response = s3.get_object(Bucket=bucket, Key=key)
        with gzip.GzipFile(fileobj=response['Body']) as log_file:
            # Parse TSV logs into a dictionary
            reader = csv.DictReader(log_file, delimiter='\t')
            for row in reader:
                # Filter only HLS media requests
                if not row['cs-uri-stem'].endswith(('.m3u8', '.ts')):
                    continue
                
                # Extract core data points
                user_id = row.get('x-edge-request-header-X-User-ID')
                if not user_id:
                    continue  # Skip unauthenticated requests
                
                # Extract video ID from path (adjust based on your file structure)
                path_parts = row['cs-uri-stem'].split('/')
                video_id = path_parts[2] if len(path_parts) >=3 else 'unknown'
                
                # Convert bytes to MB
                bytes_sent = int(row['sc-bytes'])
                mb_consumed = bytes_sent / (1024 * 1024)

                # Update MongoDB: upsert user-video entry, increment total MB
                consumption_coll.update_one(
                    {'user_id': user_id, 'video_id': video_id},
                    {
                        '$inc': {'total_mb': round(mb_consumed, 2)},
                        '$setOnInsert': {'created_at': datetime.utcnow()},
                        '$set': {'last_accessed': datetime.utcnow()}
                    },
                    upsert=True
                )
    except Exception as e:
        print(f"Error processing log file {key}: {str(e)}")
        raise

Key notes for this Lambda:

  • Store your MongoDB URI in Lambda environment variables (never hardcode credentials).
  • Set a sufficient timeout (e.g., 5 minutes) to handle large log files.
  • Add logic to track processed log files (e.g., store ETags in a DynamoDB table) to avoid duplicate processing if CloudFront retries log uploads.
Step 3: Optimize for High Traffic (Optional)

If you're dealing with massive log volumes, Lambda alone might not scale efficiently. For this scenario:

  • Use Kinesis Data Firehose to batch CloudFront logs and load them into Amazon Athena.
  • Write Athena SQL queries to pre-aggregate daily/weekly per-user-per-video consumption.
  • Trigger a Lambda function periodically to run these queries and sync aggregated results to MongoDB.

This reduces the number of direct writes to MongoDB and makes large-scale analysis easier.

Step 4: Validate & Test
  • Send test streaming requests from your mobile app with a test user ID.
  • Check that CloudFront logs include the X-User-ID header and correct file paths.
  • Verify that Lambda processes the logs and updates MongoDB with accurate MB counts.
  • Test edge cases (e.g., partial video playback, repeated requests) to ensure consumption totals are accurate.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 16:57:42