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

PyFlink 1.19如何跨AWS账户读取Kinesis流?

要实现PyFlink 1.19读取另一个AWS账户下的Kinesis流,核心是正确配置流标识与凭证参数,以下是具体步骤及问题排查:

核心配置要点

  1. 构造函数仅传流名称:FlinkKinesisConsumer的构造函数第一个参数必须是流的名称(符合[a-zA-Z0-9_.-]+正则),不能直接传ARN,这是你第一次报错的原因。
  2. 通过配置指定ARN:在consumer_config中添加stream.arn参数,指定目标流的完整ARN,这是让客户端定位到跨账户流的关键。
  3. 确保凭证权限有效:运行代码的IAM角色必须具备访问目标流的权限,且目标账户的Kinesis流资源策略已授权该角色访问。

完整代码示例

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.connectors import FlinkKinesisConsumer
from pyflink.common.serialization import JsonRowDeserializationSchema
from pyflink.common.types import Row

# 初始化执行环境
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)

# 配置目标流信息
target_stream_name = "你的目标流名称"
target_stream_arn = "arn:aws:kinesis:eu-west-1:目标账户ID:stream/你的目标流名称"
target_region = "eu-west-1"

# 定义JSON反序列化Schema(根据你的数据结构调整)
deserialization_schema = JsonRowDeserializationSchema.builder()
    .type_info(Row([("user_id", str), ("event_time", str), ("event_type", str)]))
    .json_timestamp_format_standard("ISO-8601")
    .build()

# 消费者配置
consumer_config = {
    "aws.region": target_region,
    "stream.arn": target_stream_arn,
    # 使用默认凭证链(适用于EC2/EKS/ECS等已绑定IAM角色的环境)
    "aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
    # 如果需要主动切换到目标账户的角色,替换为以下STS配置:
    # "aws.credentials.provider": "com.amazonaws.auth.STSAssumeRoleSessionCredentialsProvider",
    # "sts.role.arn": "arn:aws:iam::目标账户ID:role/允许访问的角色名",
    # "sts.role.session.name": "flink-kinesis-consumer-session"
}

# 创建Kinesis消费者
kinesis_consumer = FlinkKinesisConsumer(
    target_stream_name,
    deserialization_schema,
    consumer_config
)

# 添加数据源并执行
ds = env.add_source(kinesis_consumer)
ds.print()

env.execute("Cross-Account Kinesis Consumer Job")

错误排查

你第二次遇到的ResourceNotFoundException,大概率是以下原因:

  • 未正确配置stream.arn,导致客户端仍在当前账户下查找同名流;
  • 运行代码的IAM角色凭证未正确加载,无法访问目标账户的资源;
  • 目标流的资源策略未正确授权当前角色(需确保策略中包含当前角色的ARN,且允许kinesis:DescribeStreamSummary、kinesis:GetShardIterator、kinesis:GetRecords等操作)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 00:02:12