【发布时间】: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