【问题标题】:Python code for AWS Cloudwatch metrics & Kinesis firehose dynamic partitioningAWS Cloudwatch 指标和 Kinesis firehose 动态分区的 Python 代码
【发布时间】:2022-02-10 07:37:52
【问题描述】:

我正在使用下面的 Lambda (python) 代码来解码 cloudwatch 指标流,然后根据流数据中的 metric_name 字段创建动态分区。但是正如我在代码中编写的那样,如果流文件具有 3 种类型的 metric_name 数据,则它只选择第一个 metric 名称并基于该名称创建分区并将所有 3 metric_name 放在同一分区中。这不是预期的输出。任何人都可以帮忙处理这段代码吗?

我从 AWS 文档中了解了有关此代码的想法,并根据我的流数据对其进行了更改。(参考:https://docs.aws.amazon.com/firehose/latest/dev/dynamic-partitioning.html

代码:

from __future__ import print_function
import base64
import json
import datetime
 
# Signature for all Lambda functions that user must implement
def lambda_handler(firehose_records_input, context):
    print("Received records for processing from DeliveryStream: " + firehose_records_input['deliveryStreamArn']
          + ", Region: " + firehose_records_input['region']
          + ", and InvocationId: " + firehose_records_input['invocationId'])
 
    # Create return value.
    firehose_records_output = {'records': []}
 
    # Create result object.
    # Go through records and process them
 
    for firehose_record_input in firehose_records_input['records']:
        # Get user payload
        payload = base64.b64decode(firehose_record_input['data']).decode('utf-8')
        
        metrics = list(filter(None, payload.split("\n")))
        

        if len(metrics) > 1:
            print("Length of array = {}".format(len(metrics)))
            new_payload = []
            for metric in metrics:
                new_payload.append(json.loads(metric))
            json_value = new_payload 
            print("Record after processing")
            print(json_value)
            # exit()
            print("\n")
            # Create output Firehose record and add modified payload and record ID to it.
            firehose_record_output = {}
            
            #working partition_keys
            
            partition_keys = {"metric_name": json_value[0]['metric_name'] }
            
            # Create output Firehose record and add modified payload and record ID to it.
            firehose_record_output = {'recordId': firehose_record_input['recordId'],
                                      'data': firehose_record_input['data'],
                                      'result': 'Ok',
                                      'metadata': { 'partitionKeys': partition_keys }}
     
            # Must set proper record ID
            # Add the record to the list of output records.
     
            firehose_records_output['records'].append(firehose_record_output)
     
    # At the end return processed records
    print('**************firehose_records_output**********************')
    print(firehose_records_output)
    print('****************end********************')
    
    return firehose_records_output

流数据样本:

{
    "metric_stream_name": "MyMetricStream",
    "account_id": "12345678",
    "region": "us-east-1",
    "namespace": "AWS/EC2",
    "metric_name": "DiskWriteOps",
    "dimensions": {
        "InstanceId": "i1234"
    },
    "timestamp": 1611929698000,
    "value": {
        "count": 3.0,
        "sum": 20.0,
        "max": 18.0,
        "min": 0.0
    },
    "unit": "Seconds"
},
{
    "metric_stream_name": "MyMetricStream",
    "account_id": "12345678",
    "region": "us-east-1",
    "namespace": "AWS/EC2",
    "metric_name": "DiskReadIOps",
    "dimensions": {
        "InstanceId": "i1234"
    },
    "timestamp": 1611929698000,
    "value": {
        "count": 3.0,
        "sum": 20.0,
        "max": 18.0,
        "min": 0.0
    },
    "unit": "Seconds"
},
{
    "metric_stream_name": "MyMetricStream",
    "account_id": "12345678",
    "region": "us-east-1",
    "namespace": "AWS/EC2",
    "metric_name": "CPUUtilization",
    "dimensions": {
        "InstanceId": "i1234"
    },
    "timestamp": 1611929698000,
    "value": {
        "count": 3.0,
        "sum": 20.0,
        "max": 18.0,
        "min": 0.0
    },
    "unit": "Seconds"
}

根据上述数据,应该创建 3 个分区,但它使用 1st metric_name(DiskWriteOps) 创建 1 个分区并在其中推送所有指标(因此我认为 partition_keys = {"metric_name": json_value[0]['metric_name'] }

【问题讨论】:

    标签: python amazon-web-services aws-lambda amazon-cloudwatch amazon-kinesis


    【解决方案1】:

    如果您将该记录传递给 firehose,它将被视为一条 json 记录,因此您只能看到一个分区,因为 json_value[0]['metric_name']

    在这种情况下,您需要通过 lambda 函数或 Kinesis 分析执行一些预处理,以将该 json 拆分为 3 个单独的记录,然后使用动态分区来使用 metric_name 创建具有正确名称的分区。

    动态分区代码如下所示:

     partition_keys = {"metric_name": json_value['metric_name'] }
    

    【讨论】:

    • 似乎是这样,但我在互联网上找不到任何此类参考。我在想它是否可以作为这个 lambda 函数本身的一部分来完成。不确定运动分析,但会探索
    猜你喜欢
    • 2017-10-16
    • 2019-09-29
    • 2022-09-23
    • 1970-01-01
    • 2019-03-18
    • 2019-05-07
    • 1970-01-01
    • 2016-01-14
    • 2018-03-23
    相关资源
    最近更新 更多