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

如何在Linux系统通过Amazon Kinesis Agent为Kinesis数据流添加自定义数据?

Linux版Amazon Kinesis Agent 添加自定义信封数据方案

Linux版Kinesis Agent没有Windows版的Sink Declarations直接配置项,但可以通过以下两种方案实现需求:

1. 自定义预处理脚本(推荐)

利用Kinesis Agent的processing配置项,调用外部脚本为每条记录添加自定义字段:

  • 修改Agent配置文件(默认路径/etc/aws-kinesis/agent.json),添加processing节点指定脚本路径:
{
  "cloudwatch.emitMetrics": true,
  "flows": [
    {
      "filePattern": "/var/log/app.log",
      "kinesisStream": "your-target-stream",
      "processing": [
        {
          "script": "/usr/local/bin/wrap_record.py",
          "roleArn": "arn:aws:iam::123456789012:role/KinesisAgentProcessingRole"
        }
      ]
    }
  ]
}
  • 编写预处理脚本(以Python为例),从标准输入读取原始记录,添加自定义信封后输出到标准输出:
import sys
import time
import json
from pathlib import Path

# 持久化计数器,避免Agent重启后重置
COUNTER_FILE = "/var/lib/aws-kinesis-agent/sequence_counter.txt"

# 初始化计数器
if Path(COUNTER_FILE).exists():
    with open(COUNTER_FILE, "r") as f:
        counter = int(f.read().strip())
else:
    counter = 0

for line in sys.stdin:
    counter += 1
    # 处理原始记录
    original_data = line.strip()
    # 构建带信封的记录
    wrapped_record = {
        "sequence_id": counter,
        "timestamp": int(time.time() * 1000),
        "raw_data": original_data
    }
    # 输出处理后的记录
    print(json.dumps(wrapped_record))
    # 持久化计数器
    with open(COUNTER_FILE, "w") as f:
        f.write(str(counter))
  • 给脚本添加执行权限:chmod +x /usr/local/bin/wrap_record.py
  • 确保Agent对脚本路径和计数器文件有读写权限,同时安装脚本依赖的Python环境

2. 自定义Java处理器插件

如果需要更深度的集成,可以基于Kinesis Agent的开源代码开发自定义处理器:

  • 克隆Agent源码:git clone https://github.com/awslabs/amazon-kinesis-agent.git
  • 实现com.amazonaws.kinesis.agent.processing.RecordProcessor接口,在processRecord方法中添加信封逻辑
  • 编译打包自定义插件JAR,放到Agent的类路径(默认/usr/share/aws-kinesis-agent/lib/)
  • 在配置文件中指定自定义处理器类名:
"processing": [
  {
    "class": "com.yourcompany.kinesis.CustomEnvelopeProcessor"
  }
]

验证方法

处理完成后,可通过AWS CLI命令读取流数据验证格式:

aws kinesis get-shard-iterator --stream-name your-target-stream --shard-id shardId-000000000000 --shard-iterator-type TRIM_HORIZON
# 使用返回的ShardIterator调用get-records
aws kinesis get-records --shard-iterator <your-shard-iterator>

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 12:45:08