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

AWS CloudWatch日志导入Azure Log Analytics格式及实现咨询

我刚好处理过类似的需求,把AWS Lambda的CloudWatch日志推送到Azure Log Analytics其实不难,核心就是搞懂Azure Data Collector API要求的JSON格式,再结合AWS Lambda的CloudWatch日志触发事件来处理。下面给你详细的示例和步骤:

解决方案:将AWS Lambda CloudWatch日志导入Azure Log Analytics

一、Azure Log Analytics 要求的JSON格式

Azure的Data Collector API接受单条或多条日志组成的JSON数组,每条日志可以包含自定义字段,推荐加上TimeGenerated字段(用来指定日志生成时间,否则Azure会默认使用接收时间)。

针对Lambda的CloudWatch日志,我们可以构造这样的标准格式:

[
  {
    "TimeGenerated": "2024-05-20T12:34:56Z",
    "LambdaFunctionName": "my-serverless-function",
    "LambdaRequestId": "abc123-def456-ghi789",
    "CloudWatchLogStream": "2024/05/20/[$LATEST]abcdef123456",
    "LogMessage": "INFO: User login successful",
    "LogLevel": "INFO"
  }
]
  • TimeGenerated:必须是ISO 8601格式的UTC时间,对应CloudWatch日志的原始时间戳
  • 其余字段均为自定义字段,你可以根据需求添加,比如请求参数、错误堆栈、环境标识等

二、AWS Lambda Python 实现代码

这个Lambda需要被CloudWatch日志组触发(给每个要同步的Lambda日志组都配置触发器),代码会自动解码CloudWatch的日志事件,构造符合Azure要求的JSON结构,再发送请求到Data Collector API。

import os
import json
import gzip
import base64
import requests
from datetime import datetime

# 从Lambda环境变量读取Azure配置(避免硬编码敏感信息)
AZURE_WORKSPACE_ID = os.environ.get("AZURE_WORKSPACE_ID")
AZURE_WORKSPACE_KEY = os.environ.get("AZURE_WORKSPACE_KEY")
AZURE_LOG_TYPE = "AWSLambdaCloudWatchLogs"  # 自定义日志类型,Azure中会作为表名前缀(最终表名为AWSLambdaCloudWatchLogs_CL)

def build_azure_signature(date, content_length, method, content_type, resource):
    """生成Azure Data Collector API要求的身份签名"""
    x_headers = f"x-ms-date:{date}"
    string_to_hash = f"{method}\n{content_length}\n{content_type}\n{x_headers}\n{resource}"
    bytes_to_hash = bytes(string_to_hash, encoding="utf-8")
    
    import hmac
    import hashlib
    decoded_key = base64.b64decode(AZURE_WORKSPACE_KEY)
    encoded_hash = base64.b64encode(hmac.new(decoded_key, bytes_to_hash, digestmod=hashlib.sha256).digest()).decode()
    return f"SharedKey {AZURE_WORKSPACE_ID}:{encoded_hash}"

def send_logs_to_azure(log_data):
    """将构造好的日志数据发送到Azure Log Analytics"""
    if not AZURE_WORKSPACE_ID or not AZURE_WORKSPACE_KEY:
        raise ValueError("请在Lambda环境变量中配置Azure工作区ID和密钥")
    
    date = datetime.utcnow().strftime("%a, %d %b %Y %H:%M:%S GMT")
    content_length = len(log_data)
    method = "POST"
    content_type = "application/json"
    resource = "/api/logs"
    uri = f"https://{AZURE_WORKSPACE_ID}.ods.opinsights.azure.com{resource}?api-version=2016-04-01"
    
    signature = build_azure_signature(date, content_length, method, content_type, resource)
    
    headers = {
        "content-type": content_type,
        "Authorization": signature,
        "Log-Type": AZURE_LOG_TYPE,
        "x-ms-date": date
    }
    
    response = requests.post(uri, data=log_data, headers=headers)
    if response.status_code not in [200, 202]:
        raise Exception(f"推送日志到Azure失败:{response.status_code} - {response.text}")

def lambda_handler(event, context):
    # 解码CloudWatch传递的压缩日志数据
    compressed_payload = base64.b64decode(event['awslogs']['data'])
    payload = gzip.decompress(compressed_payload).decode('utf-8')
    log_events = json.loads(payload)
    
    # 构造符合Azure要求的日志数组
    azure_log_entries = []
    for log_event in log_events['logEvents']:
        # 转换CloudWatch毫秒级时间戳为ISO格式UTC时间
        log_timestamp = datetime.utcfromtimestamp(log_event['timestamp'] / 1000).isoformat() + "Z"
        # 尝试从日志内容中提取级别(适配常见的日志格式)
        log_level = "UNKNOWN"
        if log_event['message'].startswith(('INFO:', 'WARN:', 'ERROR:', 'DEBUG:')):
            log_level = log_event['message'].split(':')[0]
        
        azure_log_entries.append({
            "TimeGenerated": log_timestamp,
            "LambdaFunctionName": log_events['logStream'].split('/')[2].split(']')[0],
            "LambdaRequestId": context.aws_request_id,
            "CloudWatchLogGroup": log_events['logGroup'],
            "CloudWatchLogStream": log_events['logStream'],
            "LogMessage": log_event['message'].strip(),
            "LogLevel": log_level
        })
    
    # 批量发送日志到Azure
    if azure_log_entries:
        log_json = json.dumps(azure_log_entries)
        send_logs_to_azure(log_json)
    
    return {
        'statusCode': 200,
        'body': json.dumps(f"成功处理{len(azure_log_entries)}条日志事件")
    }

三、配置步骤

1. Azure端配置

  • 打开你的Log Analytics工作区,进入代理管理 -> 密钥,记录下工作区ID和主键/次键
  • 无需提前创建日志表,Azure会自动根据AZURE_LOG_TYPE创建对应的数据表(格式为[自定义类型]_CL)

2. AWS端配置

  • 创建Python Lambda函数,粘贴上述代码
  • 在Lambda的环境变量中添加AZURE_WORKSPACE_ID和AZURE_WORKSPACE_KEY,填入从Azure获取的值
  • 给Lambda添加CloudWatch日志触发器:针对每个需要同步的Lambda日志组,配置触发器,可按需设置日志过滤规则(比如仅同步ERROR级别日志)
  • 配置Lambda权限:确保角色拥有logs:DescribeLogStreams和logs:GetLogEvents权限(触发器自动创建的角色会默认包含),同时允许Lambda出站访问公网(若在VPC内,需配置NAT网关)

四、注意事项

  • 日志大小限制:Azure Data Collector API单请求最大支持1MB,若CloudWatch单次触发的日志量过大,需添加日志拆分逻辑
  • 可靠性优化:建议添加重试机制(比如使用tenacity库),避免网络波动导致日志丢失
  • 日志解析:如果你的Lambda日志是JSON格式,可以在代码中解析LogMessage,将嵌套字段单独提取到Azure日志中,方便后续查询分析
  • 成本评估:Azure Log Analytics按数据 ingestion 量收费,AWS Lambda和CloudWatch也有对应费用,需提前评估成本

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:10:23