亚马逊AWS官方博客

Kafka → Amazon Redshift 数据摄取 · 三方案对比与部署手册

摘要:三种方案:MSK 托管 + Redshift 物化视图(推荐)· 自建 Kafka(TLS)直摄 · Kafka→S3 Tables→Redshift 查询


一、概览

本文档汇总三种把 Kafka 数据接入 Amazon Redshift 的方案,均围绕近实时分析场景:

方案一:MSK 托管 + Redshift 物化视图(推荐)

托管 Kafka(MSK Serverless),IAM 认证,Redshift Streaming Ingestion + 自动刷新物化视图直接消费,秒级时延、免运维 broker。综合最省心,首选方案。

方案二:自建 Kafka(TLS)直摄

EC2 自建社区版 Kafka,公有可信证书 TLS,Redshift 直接消费,成本可控、可定制。

方案三:Kafka → S3 Tables → Redshift

先落地为 S3 Tables(Iceberg),Redshift 外部表直接查询,存算分离、可被多引擎共享。

适用场景:本文面向大规模数据采集的近实时入仓场景 —— 先用 Kafka 承接高吞吐(削峰 / 解耦 / 容错 / 重放),再由 Redshift 直接消费 Kafka 入仓分析。围绕这一需求,本文从 MSK 托管、自建 Kafka、S3 Tables 三条路径给出对比、选型与可落地的部署步骤,示例值均已脱敏,可直接套用到同类场景。

二、方案对比

维度 方案一:MSK 直摄 (推荐) 方案二:自建 Kafka 直摄 方案三:S3 Tables
Kafka 运维 全托管,无需管 broker 自运维 EC2/补丁/扩缩容 取决于上游(自建或 MSK)
认证方式 IAM(9098) TLS / mTLS(9094),需公有可信证书 由写入端决定
Redshift 接入 Streaming Ingestion(外部 schema + 物化视图) Streaming Ingestion(外部 schema + 物化视图) 外部表查询 Iceberg(Spectrum / 湖仓)
时延 秒级 秒级 分钟级(取决于落盘/compaction)
数据是否落 S3 否(直入 Redshift) 否(直入 Redshift) 是,数据湖可多引擎共享
证书/网络复杂度 低(IAM) 较高(需公有 CA 证书 + DNS) 中
成本要点 MSK Serverless 集群固定费(约 $0.75/小时) EC2 + NAT + 可导出证书,相对低 S3 存储 + 请求,存算分离最省算力
适用场景 要快速上线、不想运维 Kafka 已有自建 Kafka、需定制或控成本 需要数据湖、多引擎共享、冷热分层
选型建议(推荐 → 方案一):MSK 托管 + Redshift 物化视图是本手册的首选推荐方案 —— Kafka 全托管免运维、IAM 认证网络简单、Streaming Ingestion 配合自动刷新物化视图即可实现秒级近实时入仓,人力成本最低、上线最快,适合绝大多数团队。
其余场景:已有自建 Kafka 或对基础设施成本极度敏感且具备专职运维 → 方案二;需要把流数据沉淀成数据湖、供 Redshift/Athena/EMR 等多引擎共享查询 → 方案三。三者可组合(如 方案一/二 + 方案三 同时落 Redshift 与数据湖)。

三、成本对比

下表为 us-east-1 的粗略示例月估算,用于方案横向比较;实际以官方定价与真实用量为准,建议用 AWS Pricing Calculator 精算。Redshift 计算成本三方案共用(按 RPU 计费),故单列。

3.1 基础设施成本(示例月估算)

成本项 方案一 MSK 方案二 自建 Kafka 方案三 S3 Tables
Kafka 层 MSK Serverless 集群固定费≈\$0.75/小时(≈\$548/月)+ 分区/吞吐按量 EC2 t3.medium≈\$30/月+ gp3 30GB ≈\$2.4/月 上游写入组件(Connector/Firehose/Flink,视选型)
数据落地 / 中转 无(直入 Redshift) 无(直入 Redshift) S3 Tables 存储 + 请求 + compaction(按量,数据量驱动)
网络 同 VPC,少量 NAT 网关≈\$32/月+ 流量按量 S3 请求 / 数据传输(按量)
证书 无(IAM 认证) ACM 可导出公有证书(签发/续期收费) 视写入端而定
小计(不含 Redshift) 较高≈\$550+/月起 较低≈\$65~70/月+ 证书费 中,随数据量线性增长(存算分离,最省算力)
Redshift 计算(共用) Redshift Serverless 按 RPU-小时计费(空闲自动暂停不计算);方案三查询数据湖为 Spectrum/湖仓扫描,按扫描量 / RPU 计。

3.2 人员维护成本

方案 主要运维内容 维护成本
方案一:MSK 无需管 broker(全托管);仅维护 IAM 角色、物化视图刷新与监控。无补丁/扩缩容负担。 低
方案二:自建 Kafka 自运维 broker:OS/Kafka 打补丁、扩缩容与高可用、证书 198 天续期、监控告警、故障恢复、容量规划。 高
方案三:S3 Tables 维护上游写入管道(Connector/Firehose/Flink);S3 Tables 的 compaction/元数据由托管服务处理,表结构与分区需治理。 中

3.3 综合(基础设施 + 人力)

3.3.1 方案一:MSK(推荐)

  • 基础设施:较高
  • 人力:低
  • 花钱买省心,快速上线、免运维,物化视图自动刷新。综合最优,首选。

3.3.2 方案二:自建 Kafka

  • 基础设施:低
  • 人力:高
  • 机器便宜,但需专人运维,隐性人力成本高。

3.3.3 方案三:S3 Tables

  • 基础设施:中(按量)
  • 人力:中
  • 存算分离最省算力,适合数据湖多引擎共享。
结论(推荐方案一):综合基础设施与人力成本,MSK 托管 + Redshift 物化视图(方案一)为推荐方案 —— 虽然 MSK Serverless 有固定费,但免去 broker 运维、补丁、证书续期等隐性人力,秒级入仓且上线最快,适合绝大多数团队。仅在有专职运维团队且对机器成本极度敏感时考虑方案二;需要数据湖沉淀与多引擎共享时选方案三。以上数字为粗略示例,请以 AWS Pricing Calculator 精算为准。

四、物化视图性能与数据量考量

一句话结论:Redshift 流式摄取没有”数据量上限”的硬性限制 —— 每次刷新可摄取数百 MB/秒,官方定位就是高吞吐近实时入仓。但大数据量下有几处性能要点必须注意,否则会出现刷新滞后、延迟升高甚至丢数。手册示例(单分区、100 条)是功能演示,生产大流量需按下表调整。

4.1 大数据量下的关键要点

维度 要点 建议
并行度 = 分区数 刷新时 Redshift 把每个 Kafka 分区分配给一个计算 slice 并行消费。topic 只有 1 个分区就只有 1 路并行,是大流量最常见瓶颈。 生产按吞吐把 topic 分区数提高到数十个,使其 ≥ Redshift slice 数以充分并行。示例中的 num_partitions=1 仅用于演示。
MV 定义要”薄” 用 JSON_EXTRACT_PATH_TEXT 每取一列就重解析一次 JSON(10 列 = 解析 10 次),延迟显著升高。 沿用手册做法:JSON_PARSE 落成 SUPER,落地后再用 PartiQL 解析;保持流式 MV 尽量简单,复杂转换 / 清洗放到下游。
不支持 JOIN 流式 MV 不支持 JOIN,也不宜在其上做复杂聚合。 先建”薄”流式 MV,再建第二层 MV / 表做 JOIN、聚合与建模。
刷新须能增量维护 MV 必须可增量刷新;若刷新长期滞后、超过 Kafka 保留期(MSK 默认 7 天),超期数据无法补齐会丢数。 为大流量预留足够算力让刷新跟上;适当拉长 topic 保留期作兜底;持续监控滞后。
单条记录 ≤ 16 MiB 单条消息超过 16 MiB(VARBYTE 上限)会被跳过并记入 SYS_STREAM_SCAN_ERRORS(刷新本身仍成功)。 控制单条消息体积;超大字段拆分或改走对象存储引用。
算力争用 自动刷新与用户查询共用计算资源,大流量 + 高并发查询会互相争用。 为 Serverless 设足够 RPU / 为预置集群设足够节点;必要时错峰或隔离工作负载。
一个 topic 一个 MV 同一 topic 建多个流式 MV 会给每个分区建多个消费者,导致限流 / 超吞吐、重复计费。 每个 topic 只建一个流式 MV,下游再分流。
监控:用 SYS_STREAM_SCAN_STATES 看每次刷新的滞后 / 吞吐,用 SYS_STREAM_SCAN_ERRORS 看被跳过的错误记录。大流量上线前建议按峰值吞吐做压测,确认刷新能持续追平流。

五、方案一:MSK 托管 + Redshift 物化视图 → 直摄 · Step by Step(推荐)

把 Amazon MSK(托管 Kafka)topic 消息,通过 Redshift Streaming Ingestion + 自动刷新物化视图近实时摄取到 Redshift Serverless,全程不经过 S3。本方案为推荐方案:免运维、秒级时延、网络与认证最简单。

环境占位符(请替换为你的实际值):账号 <ACCOUNT_ID> · VPC vpc-xxxxxxxx · 安全组 sg-xxxxxxxx · 生产者 EC2 i-xxxxxxxx · namespace <namespace> · workgroup <workgroup> · 库 dev · 集群 my-msk · topic orders-topic

Step 1 · 创建 MSK Serverless 集群

  1. MSK 控制台 →Create cluster→Custom create(勿用 Quick create,否则无法指定 VPC)
  2. Cluster name:my-msk;Cluster type:Serverless
  3. 网络:选你的 VPC vpc-xxxxxxxx,至少 2 个不同 AZ 的子网
  4. 安全组:选与 Redshift 可互通的安全组 sg-xxxxxxxx
  5. 认证:IAM;创建(约十几分钟变 Active)

获取 bootstrap 地址(IAM 端口 9098):

aws kafka list-clusters-v2 --region us-east-1
aws kafka get-bootstrap-brokers --region us-east-1 --cluster-arn <ClusterArn>
# 输出示例(已脱敏):
# "BootstrapBrokerStringSaslIam": "boot-xxxxxxxx.c2.kafka-serverless.us-east-1.amazonaws.com:9098"

Step 2 · 确认网络连通

IAM 认证使用 9098。确认安全组入站允许来自 Redshift/EC2 的流量,并确认 Redshift 已开 Enhanced VPC routing:

aws redshift-serverless list-workgroups --region us-east-1 \
  --query 'workgroups[].{name:workgroupName,evr:enhancedVpcRouting}'

# 从 EC2 实测 9098 连通(SSM 下发)
aws ssm send-command --region us-east-1 --instance-ids i-xxxxxxxx \
  --document-name AWS-RunShellScript \
  --parameters 'commands=["timeout 8 bash -c \"cat < /dev/null > /dev/tcp/<bootstrap 主机>/9098\" && echo TCP_OK || echo TCP_FAIL"]'

Step 3 · 为 Redshift 创建 IAM 角色(只读 MSK)

信任策略 redshift-trust-policy.json

{
  "Version": "2012-10-17",
  "Statement": [{
    "Effect": "Allow",
    "Principal": { "Service": "redshift.amazonaws.com" },
    "Action": "sts:AssumeRole"
  }]
}

权限策略 redshift-msk-policy.json(只读)

{
  "Version": "2012-10-17",
  "Statement": [
    { "Sid":"Connect","Effect":"Allow",
      "Action":["kafka-cluster:Connect","kafka-cluster:DescribeCluster"],
      "Resource":"arn:aws:kafka:us-east-1:<ACCOUNT_ID>:cluster/my-msk/*" },
    { "Sid":"ReadTopic","Effect":"Allow",
      "Action":["kafka-cluster:ReadData","kafka-cluster:DescribeTopic"],
      "Resource":"arn:aws:kafka:us-east-1:<ACCOUNT_ID>:topic/my-msk/*" },
    { "Sid":"Group","Effect":"Allow",
      "Action":["kafka-cluster:AlterGroup","kafka-cluster:DescribeGroup"],
      "Resource":"arn:aws:kafka:us-east-1:<ACCOUNT_ID>:group/my-msk/*" }
  ]
}

创建并附加,挂到 namespace

aws iam create-role --role-name RedshiftMSKStreamingRole \
  --assume-role-policy-document file://redshift-trust-policy.json
aws iam create-policy --policy-name RedshiftMSKStreamingPolicy \
  --policy-document file://redshift-msk-policy.json
aws iam attach-role-policy --role-name RedshiftMSKStreamingRole \
  --policy-arn arn:aws:iam::<ACCOUNT_ID>:policy/RedshiftMSKStreamingPolicy

# 挂到 namespace(保留原有默认角色)
aws redshift-serverless update-namespace --region us-east-1 \
  --namespace-name <namespace> \
  --default-iam-role-arn <原有默认角色 ARN> \
  --iam-roles <原有默认角色 ARN> arn:aws:iam::<ACCOUNT_ID>:role/RedshiftMSKStreamingRole

Step 4 · 为生产者 EC2 授予 MSK 写权限

{
  "Version":"2012-10-17",
  "Statement":[
    {"Sid":"Cluster","Effect":"Allow",
     "Action":["kafka-cluster:Connect","kafka-cluster:AlterCluster","kafka-cluster:DescribeCluster"],
     "Resource":"arn:aws:kafka:us-east-1:<ACCOUNT_ID>:cluster/my-msk/*"},
    {"Sid":"Topic","Effect":"Allow",
     "Action":["kafka-cluster:*Topic*","kafka-cluster:WriteData","kafka-cluster:ReadData"],
     "Resource":"arn:aws:kafka:us-east-1:<ACCOUNT_ID>:topic/my-msk/*"},
    {"Sid":"Group","Effect":"Allow",
     "Action":["kafka-cluster:AlterGroup","kafka-cluster:DescribeGroup"],
     "Resource":"arn:aws:kafka:us-east-1:<ACCOUNT_ID>:group/my-msk/*"}
  ]
}

aws iam create-policy --policy-name EC2MSKProducerPolicy \
  --policy-document file://ec2-msk-producer-policy.json
aws iam attach-role-policy --role-name <EC2 实例角色> \
  --policy-arn arn:aws:iam::<ACCOUNT_ID>:policy/EC2MSKProducerPolicy

Step 5 · 在 EC2 准备 Python 生产环境

aws ssm send-command --region us-east-1 --instance-ids i-xxxxxxxx \
  --document-name AWS-RunShellScript \
  --parameters 'commands=[
    "dnf install -y python3-pip",
    "python3 -m pip install \"kafka-python==2.0.2\" aws-msk-iam-sasl-signer-python",
    "python3 -c \"import kafka; print(kafka.__version__)\""
  ]'
关键坑:kafka-python 必须固定 2.0.2。默认可能装到 3.x(异步重写版),其 OAUTHBEARER 与 AWS 签名器不兼容,会报 KafkaTimeoutError: Unable to bootstrap(即便网络/IAM 都正常)。

Step 6 · 生产测试数据(订单宽表 100 条)

/opt/order_producer.py:清空并重建 topic → 生产 100 条订单大宽表 JSON(37 字段,含订单/账户/收货地址/商品/物流)。

import json, time, random, datetime
from kafka import KafkaProducer
from kafka.admin import KafkaAdminClient, NewTopic
from kafka.errors import TopicAlreadyExistsError, UnknownTopicOrPartitionError
from aws_msk_iam_sasl_signer import MSKAuthTokenProvider

BOOTSTRAP = 'boot-xxxxxxxx.c2.kafka-serverless.us-east-1.amazonaws.com:9098'
REGION = 'us-east-1'
TOPIC = 'orders-topic'

class TP:
    def token(self):
        t, _ = MSKAuthTokenProvider.generate_auth_token(REGION); return t

common = dict(bootstrap_servers=BOOTSTRAP, security_protocol='SASL_SSL',
              sasl_mechanism='OAUTHBEARER', sasl_oauth_token_provider=TP())

# 清空:删除 topic 再重建(offset 归零)
admin = KafkaAdminClient(**common, client_id='admin1')
try: admin.delete_topics([TOPIC])
except UnknownTopicOrPartitionError: pass
for _ in range(30):
    time.sleep(3)
    if TOPIC not in admin.list_topics(): break
for _ in range(20):
    try:
        admin.create_topics([NewTopic(name=TOPIC, num_partitions=1, replication_factor=3)]); break
    except TopicAlreadyExistsError: time.sleep(3)
admin.close()

surnames=['张','王','李','赵','刘','陈','杨','黄','周','吴']; names=['伟','芳','娜','强','磊','军','洋','勇','艳','杰']
provinces=[('广东省','深圳市','南山区'),('北京市','北京市','朝阳区'),('浙江省','杭州市','西湖区'),('四川省','成都市','武侯区')]
streets=['科技园路','人民大道','解放路','中山路']; levels=['普通会员','银卡会员','金卡会员','钻石会员']
statuses=['待付款','已付款','已发货','已完成','已取消']; channels=['APP','小程序','Web','门店']; pays=['支付宝','微信支付','银行卡','货到付款']
products=[('P1001','小米手机 14','手机数码','小米',3999),('P1004','戴森吹风机','家用电器','Dyson',2999),('P1007','三只松鼠坚果','食品生鲜','三只松鼠',89)]
express=['顺丰速运','圆通速递','中通快递','京东物流']; warehouses=['华南仓','华东仓','华北仓','西南仓']

producer=KafkaProducer(**common, value_serializer=lambda v: json.dumps(v, ensure_ascii=False).encode('utf-8'))
base=datetime.datetime(2026,9,1)
for i in range(1,101):
    cust=random.choice(surnames)+random.choice(names); recv=random.choice(surnames)+random.choice(names)
    prov,city,dist=random.choice(provinces); p=random.choice(products); qty=random.randint(1,5); unit=p[4]
    subtotal=unit*qty; discount=random.choice([0,0,0,10,50,100]); freight=random.choice([0,0,8,12,15]); actual=subtotal-discount+freight
    ctime=base+datetime.timedelta(days=random.randint(0,9),hours=random.randint(0,23)); status=random.choice(statuses)
    rec={'order_id':'ORD%08d'%i,'order_no':'SY%s%04d'%(ctime.strftime('%Y%m%d'),i),'order_status':status,
      'order_channel':random.choice(channels),'created_at':ctime.isoformat(),
      'paid_at':(ctime+datetime.timedelta(minutes=30)).isoformat() if status!='待付款' else None,
      'currency':'CNY','total_amount':subtotal,'discount_amount':discount,'freight':freight,'actual_amount':actual,
      'payment_method':random.choice(pays),'account_id':'ACC%06d'%random.randint(1,60),'customer_name':cust,
      'gender':random.choice(['男','女']),'phone':'138%08d'%random.randint(0,99999999),'email':'user%d@example.com'%random.randint(1000,9999),
      'member_level':random.choice(levels),'register_date':'2025-01-15','points':random.randint(0,50000),
      'receiver_name':recv,'receiver_phone':'139%08d'%random.randint(0,99999999),'country':'中国','province':prov,'city':city,'district':dist,
      'street_address':random.choice(streets)+'%d 号'%random.randint(1,999),'postal_code':'%06d'%random.randint(100000,999999),
      'product_id':p[0],'product_name':p[1],'category':p[2],'brand':p[3],'sku':'SKU%s%03d'%(p[0],random.randint(1,999)),
      'quantity':qty,'unit_price':unit,'subtotal':subtotal,'warehouse':random.choice(warehouses),
      'shipping_company':random.choice(express),'tracking_no':'SF%013d'%random.randint(0,9999999999999)}
    producer.send(TOPIC, rec)
producer.flush(); producer.close(); print('PRODUCE_DONE 100 orders')

运行:

aws ssm send-command --region us-east-1 --instance-ids i-xxxxxxxx \
  --document-name AWS-RunShellScript \
  --parameters 'commands=["cd /opt && python3 order_producer.py"]'
# 预期: PRODUCE_DONE 100 orders

Step 7 · Redshift 建外部 Schema 与物化视图

Query Editor v2 连 <workgroup> 的 dev 库:

-- 外部 schema(连 MSK,IAM 认证)
CREATE EXTERNAL SCHEMA msk_schema
FROM KAFKA
IAM_ROLE 'arn:aws:iam::<ACCOUNT_ID>:role/RedshiftMSKStreamingRole'
AUTHENTICATION IAM
URI 'boot-xxxxxxxx.c2.kafka-serverless.us-east-1.amazonaws.com:9098';

-- 物化视图(自动刷新, JSON→SUPER)
CREATE MATERIALIZED VIEW mv_orders_raw AUTO REFRESH YES AS
SELECT kafka_partition, kafka_offset, kafka_timestamp,
       JSON_PARSE(kafka_value) AS payload
FROM msk_schema."orders-topic";

-- 首刷 + 验证
REFRESH MATERIALIZED VIEW mv_orders_raw;
SELECT COUNT(*) FROM mv_orders_raw;   -- 应为 100

宽表解析与业务分析示例:

SELECT payload.order_id::varchar   AS order_id,
       payload.member_level::varchar AS member_level,
       payload.actual_amount::int   AS actual_amount,
       payload.province::varchar    AS province
FROM mv_orders_raw ORDER BY kafka_offset;

-- 按会员等级
SELECT payload.member_level::varchar AS member_level,
       COUNT(*) AS order_cnt, SUM(payload.actual_amount::int) AS total_amount
FROM mv_orders_raw GROUP BY 1 ORDER BY total_amount DESC;

Step 8 · (可选)Kafka UI 通过 SSM 访问

# EC2 上用 Docker 跑 Kafka UI(SSM 下发,IAM 认证)
docker run -d --name kafka-ui --restart unless-stopped -p 8080:8080 \
  -e KAFKA_CLUSTERS_0_NAME=my-msk \
  -e KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS=boot-xxxxxxxx.c2.kafka-serverless.us-east-1.amazonaws.com:9098 \
  -e KAFKA_CLUSTERS_0_PROPERTIES_SECURITY_PROTOCOL=SASL_SSL \
  -e KAFKA_CLUSTERS_0_PROPERTIES_SASL_MECHANISM=AWS_MSK_IAM \
  -e "KAFKA_CLUSTERS_0_PROPERTIES_SASL_JAAS_CONFIG=software.amazon.msk.auth.iam.IAMLoginModule required;" \
  -e KAFKA_CLUSTERS_0_PROPERTIES_SASL_CLIENT_CALLBACK_HANDLER_CLASS=software.amazon.msk.auth.iam.IAMClientCallbackHandler \
  provectuslabs/kafka-ui:latest

# 本地 SSM 端口转发后访问 http://localhost:8080
aws ssm start-session --region us-east-1 --target i-xxxxxxxx \
  --document-name AWS-StartPortForwardingSession \
  --parameters '{"portNumber":["8080"],"localPortNumber":["8080"]}'
Kafka UI 默认无登录认证,但仅通过 SSM 隧道可达(安全组未开 8080),风险可控;长期使用建议加认证。
成本提醒:MSK Serverless 有约 $0.75/小时 的集群固定费(≈$548/月),请纳入成本评估(详见成本对比)。

六、方案二:自建 Kafka(TLS)→ Redshift 直摄 · Step by Step

EC2 自建社区版 Apache Kafka,配公有可信证书走 TLS,Redshift Streaming Ingestion 直接消费。

关键约束:① Redshift 连自建 Kafka 必须 TLS,不支持 PLAINTEXT;② 不支持自签名证书,broker 证书须来自公有可信 CA(本文用 ACM 可导出公有证书);③ 不支持 SASL/SCRAM、SASL/PLAINTEXT;④ workgroup 需开 Enhanced VPC Routing;⑤ topic 名区分大小写,SQL 用双引号。
环境占位符:账号 <ACCOUNT_ID> · 区域 us-east-1 · VPC vpc-xxxxxxxx · 私有子网 subnet-xxxxxxxx · Kafka 安全组 sg-kafka-xxxx · Redshift 安全组 sg-redshift-xxxx · broker 域名 kafka.example.com · Route53 区 ZXXXXXXXX · workgroup <workgroup> · 库 dev · topic kafka-stream-test

Step 1 · 创建 IAM 角色(SSM 免密登录)

aws iam create-role --role-name kafka-broker-ssm-role \
  --assume-role-policy-document '{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"Service":"ec2.amazonaws.com"},"Action":"sts:AssumeRole"}]}'
aws iam attach-role-policy --role-name kafka-broker-ssm-role \
  --policy-arn arn:aws:iam::aws:policy/AmazonSSMManagedInstanceCore
aws iam create-instance-profile --instance-profile-name kafka-broker-ssm-profile
aws iam add-role-to-instance-profile --instance-profile-name kafka-broker-ssm-profile \
  --role-name kafka-broker-ssm-role

Step 2 · 创建安全组并放行端口

KAFKA_SG=$(aws ec2 create-security-group --region us-east-1 \
  --group-name kafka-broker-sg --description "Kafka broker SG" \
  --vpc-id vpc-xxxxxxxx --query GroupId --output text)

# 放行 9092(本地测试) 与 9094(TLS 给 Redshift),仅来自 Redshift 安全组
aws ec2 authorize-security-group-ingress --region us-east-1 --group-id $KAFKA_SG \
  --ip-permissions \
  "IpProtocol=tcp,FromPort=9092,ToPort=9092,UserIdGroupPairs=[{GroupId=sg-redshift-xxxx}]" \
  "IpProtocol=tcp,FromPort=9094,ToPort=9094,UserIdGroupPairs=[{GroupId=sg-redshift-xxxx}]"

Step 3 · 启动 EC2 并自动安装 Kafka(user-data)

私有子网 + NAT 出网;t3.medium;SSM 访问(无公网 IP、无密钥)。kafka-userdata.sh:

#!/bin/bash
set -xe
exec > /var/log/kafka-bootstrap.log 2>&1
KAFKA_VER=4.1.2; SCALA_VER=2.13; KAFKA_HOME=/opt/kafka
dnf -y install java-17-amazon-corretto-headless tar unzip
cd /opt
# 用快 CDN dlcdn(勿用慢的 archive.apache.org)
curl -fSL --retry 3 -o kafka.tgz "https://dlcdn.apache.org/kafka/${KAFKA_VER}/kafka_${SCALA_VER}-${KAFKA_VER}.tgz"
tar -xzf kafka.tgz && ln -sfn "/opt/kafka_${SCALA_VER}-${KAFKA_VER}" "$KAFKA_HOME"
mkdir -p /var/lib/kafka-data
TOKEN=$(curl -s -X PUT "http://169.254.169.254/latest/api/token" -H "X-aws-ec2-metadata-token-ttl-seconds: 21600")
PRIVATE_IP=$(curl -s -H "X-aws-ec2-metadata-token: $TOKEN" http://169.254.169.254/latest/meta-data/local-ipv4)
CONFIG="$KAFKA_HOME/config/server.properties"     # Kafka 4.x 路径
sed -i "s#^log.dirs=.*#log.dirs=/var/lib/kafka-data#" "$CONFIG"
sed -i "s#^advertised.listeners=.*#advertised.listeners=PLAINTEXT://${PRIVATE_IP}:9092#" "$CONFIG"
KID=$("$KAFKA_HOME/bin/kafka-storage.sh" random-uuid)
"$KAFKA_HOME/bin/kafka-storage.sh" format -t "$KID" -c "$CONFIG" --standalone   # 4.x 用 --standalone
cat >/etc/systemd/system/kafka.service <<EOF
[Unit]
Description=Apache Kafka (KRaft)
After=network-online.target
Wants=network-online.target
[Service]
Type=simple
ExecStart=$KAFKA_HOME/bin/kafka-server-start.sh $CONFIG
ExecStop=$KAFKA_HOME/bin/kafka-server-stop.sh
Restart=on-failure
LimitNOFILE=100000
[Install]
WantedBy=multi-user.target
EOF
systemctl daemon-reload && systemctl enable --now kafka

启动实例:

INSTANCE_ID=$(aws ec2 run-instances --region us-east-1 \
  --image-id ami-xxxxxxxx --instance-type t3.medium \
  --subnet-id subnet-xxxxxxxx --security-group-ids $KAFKA_SG \
  --iam-instance-profile Name=kafka-broker-ssm-profile \
  --metadata-options 'HttpEndpoint=enabled,HttpTokens=required' \
  --block-device-mappings '[{"DeviceName":"/dev/xvda","Ebs":{"VolumeSize":30,"VolumeType":"gp3","Encrypted":true,"DeleteOnTermination":true}}]' \
  --user-data file://kafka-userdata.sh \
  --tag-specifications 'ResourceType=instance,Tags=[{Key=Name,Value=kafka-broker}]' \
  --query 'Instances[0].InstanceId' --output text)
踩坑:archive.apache.org 下载极慢,改用 dlcdn.apache.org(仅存当前版本,如 4.1.x)。Kafka 4.x 与 3.x 差异:配置路径 config/server.properties、格式化用 --standalone。

Step 4 · 验证 Kafka 已启动(SSM)

aws ssm start-session --target i-xxxxxxxx --region us-east-1
# 进入后:
sudo systemctl is-active kafka                  # 期望 active
sudo /opt/kafka/bin/kafka-topics.sh --version   # 期望 4.1.2
# 本机 9092 收发测试
B=/opt/kafka/bin; S=localhost:9092
sudo $B/kafka-topics.sh --bootstrap-server $S --create --topic kafka-stream-test --partitions 1 --replication-factor 1
echo '{"id":1,"name":"alice","city":"Beijing"}' | sudo $B/kafka-console-producer.sh --bootstrap-server $S --topic kafka-stream-test
sudo $B/kafka-console-consumer.sh --bootstrap-server $S --topic kafka-stream-test --from-beginning --max-messages 1

Step 5 · 申请 ACM 可导出公有证书

ACM 自 2025-06-17 起支持“可导出公有证书”:可导出私钥装到自建服务器,由 Amazon Trust Services 签发、公有可信。必须申请时就选 Export=ENABLED(签发后不可改);有效期 198 天(ACM 于到期前 45 天自动续期);签发/续期收费。
CERT_ARN=$(aws acm request-certificate --region us-east-1 \
  --domain-name kafka.example.com --key-algorithm RSA_2048 \
  --validation-method DNS --options Export=ENABLED \
  --query CertificateArn --output text)

# 取 DNS 验证记录
aws acm describe-certificate --region us-east-1 --certificate-arn $CERT_ARN \
  --query 'Certificate.DomainValidationOptions[0].ResourceRecord'

Step 6 · 写入 DNS(验证 CNAME + broker A 记录)

aws route53 change-resource-record-sets --hosted-zone-id ZXXXXXXXX \
  --change-batch '{"Changes":[
    {"Action":"UPSERT","ResourceRecordSet":{"Name":"<验证 CNAME 名>","Type":"CNAME","TTL":300,"ResourceRecords":[{"Value":"<验证 CNAME 值>"}]}},
    {"Action":"UPSERT","ResourceRecordSet":{"Name":"kafka.example.com.","Type":"A","TTL":60,"ResourceRecords":[{"Value":"<broker 私有 IP>"}]}}
  ]}'

# 等待 ISSUED
aws acm describe-certificate --region us-east-1 --certificate-arn $CERT_ARN --query 'Certificate.Status'

Step 7 · 给实例角色授予证书导出权限

aws iam put-role-policy --role-name kafka-broker-ssm-role \
  --policy-name acm-export-kafka-cert \
  --policy-document '{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"acm:ExportCertificate","Resource":"'$CERT_ARN'"}]}'

Step 8 · broker 上配置 9094 TLS 监听(私钥不外传)

SSM 登录 broker 后 sudo bash 执行(替换 ARN、DOMAIN):

ARN=<CERT_ARN>; REGION=us-east-1; DOMAIN=kafka.example.com
SSLDIR=/opt/kafka/ssl; CONFIG=/opt/kafka/config/server.properties

# 安装 aws cli v2(AL2023 默认无)
dnf install -y unzip >/dev/null 2>&1 || true
if ! command -v aws >/dev/null 2>&1; then
  curl -s https://awscli.amazonaws.com/awscli-exe-linux-x86_64.zip -o /tmp/awscliv2.zip
  cd /tmp && unzip -q -o awscliv2.zip && ./aws/install --update >/dev/null 2>&1
fi
export PATH=/usr/local/bin:$PATH

# 从 ACM 导出证书 → 生成 PKCS12 keystore
mkdir -p $SSLDIR && cd $SSLDIR
PP=$(openssl rand -base64 24); printf '%s' "$PP" > pp.txt
KSPASS=$(openssl rand -base64 24)
aws acm export-certificate --certificate-arn "$ARN" --passphrase fileb://pp.txt --region "$REGION" > export.json
python3 -c "import json;d=json.load(open('export.json'));open('cert.pem','w').write(d['Certificate']);open('chain.pem','w').write(d.get('CertificateChain','') or '');open('key.enc.pem','w').write(d['PrivateKey'])"
openssl pkey -in key.enc.pem -passin file:pp.txt -out key.pem
cat cert.pem chain.pem > fullchain.pem
openssl pkcs12 -export -in fullchain.pem -inkey key.pem -out kafka.keystore.p12 -passout pass:"$KSPASS" -name kafka
chmod 600 kafka.keystore.p12
rm -f export.json key.enc.pem pp.txt cert.pem chain.pem fullchain.pem key.pem

# 配置监听:保留 9092 PLAINTEXT,新增 9094 SSL
TOKEN=$(curl -s -X PUT http://169.254.169.254/latest/api/token -H 'X-aws-ec2-metadata-token-ttl-seconds: 21600')
PRIVATE_IP=$(curl -s -H "X-aws-ec2-metadata-token: $TOKEN" http://169.254.169.254/latest/meta-data/local-ipv4)
sed -i "s#^listeners=.*#listeners=PLAINTEXT://:9092,CONTROLLER://:9093,SSL://:9094#" $CONFIG
sed -i "s#^advertised.listeners=.*#advertised.listeners=PLAINTEXT://${PRIVATE_IP}:9092,SSL://${DOMAIN}:9094#" $CONFIG
cat >> $CONFIG <<EOF
ssl.keystore.location=/opt/kafka/ssl/kafka.keystore.p12
ssl.keystore.password=${KSPASS}
ssl.key.password=${KSPASS}
ssl.keystore.type=PKCS12
ssl.client.auth=none
EOF
chmod 600 $CONFIG
systemctl restart kafka

验证 TLS(期望 Verify return code: 0 (ok)):

echo | openssl s_client -connect kafka.example.com:9094 -servername kafka.example.com 2>/dev/null \
  | grep -E 'subject=|issuer=|Verify return code'

Step 9 · Redshift 建外部 schema 与物化视图(语义前缀 selfbuilt_kafka_)

Query Editor v2 连 <workgroup> / dev。若普通库用户报 permission denied for database dev,改用管理员凭据(Data API --secret-arn)或先授权。

-- 外部 schema(none + TLS,无需 IAM_ROLE)
CREATE EXTERNAL SCHEMA selfbuilt_kafka_ext
FROM KAFKA
URI 'kafka.example.com:9094'
AUTHENTICATION none;

-- 物化视图(topic 名含连字符,必须双引号)
CREATE MATERIALIZED VIEW selfbuilt_kafka_mv AUTO REFRESH YES AS
SELECT kafka_partition, kafka_offset, kafka_timestamp, kafka_key,
       kafka_value, JSON_PARSE(kafka_value) AS data
FROM selfbuilt_kafka_ext."kafka-stream-test"
WHERE CAN_JSON_PARSE(kafka_value);

-- 刷新 + 查询
REFRESH MATERIALIZED VIEW selfbuilt_kafka_mv;
SELECT kafka_offset, data.id::int AS id, data.name::text AS name, data.city::text AS city
FROM selfbuilt_kafka_mv ORDER BY kafka_offset;

-- 监控
SELECT * FROM SYS_STREAM_SCAN_STATES ORDER BY record_time DESC;
SELECT * FROM SYS_STREAM_SCAN_ERRORS ORDER BY record_time DESC;

Data API 方式(用管理员密钥):

SECRET=<admin-secret-arn>
aws redshift-data execute-statement --region us-east-1 \
  --workgroup-name <workgroup> --database dev --secret-arn $SECRET \
  --sql "REFRESH MATERIALIZED VIEW selfbuilt_kafka_mv;"

七、方案三 · Kafka → S3 Tables → Redshift 查询

数据先由 Kafka 写入 Amazon S3 Tables(基于 Apache Iceberg 的托管表存储),Redshift 通过外部表直接查询 S3 Tables,实现存算分离、数据湖多引擎共享。此处仅给出架构与简要说明。

[图 1]

7.1 简要说明

  • 写入:Kafka 数据通过 Kafka Connect S3/Iceberg Sink、Amazon Data Firehose 或 Flink 落地为 S3 Tables(Iceberg 格式)。注:用 Firehose 写 Iceberg/S3 Tables 目标时不支持 Amazon MSK Serverless 作为源(官方限制),需改用 MSK Provisioned,或用 Kafka Connect / Flink / 自研消费程序(如 pyiceberg)写入。
  • 查询:Redshift 通过外部 schema/外部表直接查询 S3 Tables,数据不搬进 Redshift,实现存算分离。
  • 共享:同一份数据湖表可被 Athena、EMR、Spark、其他 Redshift 集群等多引擎共享,避免重复拷贝。
  • 取舍:时延为分钟级(依赖落盘与 compaction),适合数据湖 / 冷热分层 / 多引擎共享;对秒级实时更适合方案一或二。
本方案仅作架构示意,详细落地步骤可另行编写(可复用现有《Redshift × S3 Tables 操作手册》)。

Redshift 查询 S3 Tables 的外部 Schema(可直接套用)

Redshift 通过 CREATE EXTERNAL SCHEMA ... FROM DATA CATALOG 引用 S3 Tables 的 Glue 目录。须以 IAM 身份连接 Redshift(非数据库本地用户),否则报 Database not found。

-- 方式一:指定角色 ARN(标准写法)
CREATE EXTERNAL SCHEMA s3t
FROM DATA CATALOG DATABASE '<namespace>'
IAM_ROLE 'arn:aws:iam::<ACCOUNT_ID>:role/RedshiftS3TablesRole'
REGION 'us-east-1'
CATALOG_ID '<ACCOUNT_ID>:s3tablescatalog/<表桶名>';

-- 方式二:IAM 联合身份直通(推荐,交由 Lake Formation 细粒度鉴权)
CREATE EXTERNAL SCHEMA s3t
FROM DATA CATALOG DATABASE '<namespace>'
CATALOG_ID '<ACCOUNT_ID>:s3tablescatalog/<表桶名>'
IAM_ROLE 'SESSION'
CATALOG_ROLE 'SESSION';

SELECT * FROM s3t.<表名>;
网络注意:Redshift 开启增强 VPC 路由且无 NAT 时,查询 S3 Tables 需增加 Glue 接口终端节点并放行出站 443,否则查询超时。

八、FAQ / 常见排查

现象 原因 / 解决 方案
KafkaTimeoutError: Unable to bootstrap kafka-python 版本问题,必须用 2.0.2 一
Kafka 下载卡住/极慢 用 dlcdn.apache.org,勿用 archive.apache.org 二
TLS handshake 失败 用了自签名证书(不支持);或域名与证书 CN/SAN 不匹配 二
连接超时 安全组端口未放行 / 跨 VPC 路由不通 / 域名解析失败 一、二
摄取 0 条 topic 名大小写/双引号问题;或数据灌进了明文 9092 而非 9094 二
permission denied for database dev 库用户无权限,改用管理员凭据或先授权 一、二
CTAS 普通表不更新 静态快照;要自动更新用物化视图或定时 INSERT 通用
能删指定几条消息吗 不能;MSK Serverless 不支持 DeleteRecords,靠保留期过期或删除重建 topic 一
中国区能用吗 MSK Serverless / Redshift Serverless / S3 Tables 在北京、宁夏均已提供,方案一、二整体可移植;但方案三的 Firehose 写 Iceberg/S3 Tables 在中国区不可用,需改用消费程序(pyiceberg)或 MSK Provisioned + Connector/Flink。流式摄取 SQL、ACM 可导出证书建议在目标集群二次确认,详见中国区适用性 通用
大数据量下 MV 会不会慢 / 丢数 提高 topic 分区数、保持 MV 定义薄、拉长保留期、预留算力;详见物化视图性能与数据量 一、二

九、中国区(北京 cn-north-1 / 宁夏 cn-northwest-1)适用性

结论:本文示例基于 us-east-1,但三方案所依赖的核心服务 —— MSK Serverless、MSK Provisioned、Redshift Serverless、S3 Tables —— 在中国区(北京、宁夏)均已提供,方案整体可移植到中国区。中国区由光环新网(北京)/ 西云数据(宁夏)运营,服务上线节奏可能与其他区域不同,请以目标区控制台实际可见为准。
能力 / 服务 中国区可用性 确认方式 说明
Amazon Redshift / Redshift Serverless 可用 ✓ aws redshift-serverless list-workgroups 北京、宁夏均已提供,入仓目标本身不受影响。
Amazon MSK(Provisioned 预置) 可用 ✓ aws kafka list-clusters-v2 两区均可用。
Amazon MSK Serverless 可用 ✓ 控制台 /aws kafka create-cluster-v2 --serverless 方案一(推荐)可在中国区使用。如遇文档区域列表滞后,以控制台 / CLI 实际可见为准。
Amazon S3 Tables(Iceberg 托管表) 可用 ✓ aws s3tables list-table-buckets(端点 s3tables.cn-*.amazonaws.com.rproxy.govskope.ca.cn) 两区均可用;但用 Firehose 写 Iceberg/S3 Tables 目标在中国区不可用(见下一行),方案三写入需改用消费程序(如 pyiceberg)或 MSK Provisioned + Connector/Flink。
Amazon Data Firehose(Iceberg / S3 Tables 目标) 不可用 ✗ API 实测 /官方文档 官方明确:除中国区、GovCloud(US)等区域外均支持,即北京、宁夏均无 Iceberg/S3 Tables 目标。报错 Firehose with Iceberg as destination is not available in the current region。
Glue Data Catalog / Athena / Lake Formation 可用 ✓ 控制台 / CLI 方案三数据湖侧配套齐全。
Amazon S3 / ACM(基础能力) 可用 ✓ 控制台 / CLI S3、ACM 基础能力可用。
提示:各区域的服务可用性会随时间更新。规划落地前,可在目标区用 AWS 控制台或 CLI(如 aws kafka list-clusters-v2、aws s3tables list-table-buckets、aws redshift-serverless list-workgroups)快速确认所需服务是否可见。
中国区选型建议:三方案整体可移植 —— ① 托管 Kafka → MSK Serverless(推荐)或 MSK Provisioned + Redshift 物化视图;② 已自建 Kafka → 方案二可移植,重点确认 broker 公有可信证书来源(注:实测中 Amazon 托管账号会被拒绝导出公有证书,报错 Accounts managed by Amazon or AWS are disallowed from exporting public certificates,属账号类型限制而非区域限制;请用客户自有账号 + 已注册域名先验证证书导出与 9094 TLS 链路);③ 数据湖共享 → S3 Tables + Redshift 外部表查询(方案三写入不能用 Firehose,见上表)。以落地时目标区控制台实际可见服务为准。

9.1 中国区落地改写要点(基于北京 cn-north-1 实测)

改写项 us-east-1 写法 中国区写法
ARN 分区 arn:aws:... arn:aws-cn:...
服务 Endpoint *.amazonaws.com *.amazonaws.com.rproxy.govskope.ca.cn(如 s3tables---cn-north-1.amazonaws.com.rproxy.govskope.ca.cn)
IAM 信任主体 ec2.amazonaws.com ec2.amazonaws.com.rproxy.govskope.ca.cn
成本口径 美元(MSK Serverless ≈$0.75/小时) 人民币(MSK Serverless 北京 ≈¥7.914/集群·小时≈¥5,777/月,宁夏 ≈¥5.297≈¥3,867/月)
软件下载 Docker Hub / PyPI / dlcdn.apache.org 改走国内镜像或 ECR,或全私网时经 S3 分发(Java/Kafka 包)

十、参考链接

➡️ 下一步行动:

相关产品:

  • Amazon Redshift — 经济高效的数据仓库
  • Amazon MSK — 完全托管式 Apache Kafka 服务
  • Amazon S3 — 适用于 AI、分析和存档的几乎无限的安全对象存储
  • Amazon IAM — 身份管理和访问权限
  • Amazon EC2 — 安全且可调整大小的计算容量

相关文章:

*前述特定亚马逊云科技生成式人工智能相关的服务目前在亚马逊云科技海外区域可用。亚马逊云科技中国区域相关云服务由西云数据和光环新网运营,具体信息以中国区域官网为准。

本篇作者

张振华

亚马逊云科技解决方案架构师,曾在携程、爱乐奇等互联网公司担任核心技术岗位,积累了丰富的系统架构设计经验。在AWS,主要负责云计算方案的架构设计与咨询,并协助企业在生成式 AI 方向的探索与落地实践。在 Edge、Serverless、容器化、微服务架构、云原生 DevOps、AI Agent 及 GenAI 企业级应用等领域拥有丰富的实战经验,致力于帮助企业构建云原生架构下的 AI 创新方案,加速 AI 业务落地,实现合规、安全的智能体开发与部署。

赵阳阳

亚马逊云科技解决方案架构师,有超过 10 年的研发及架构设计经验。目前致力于推广 AWS 的技术和各种解决方案。


AWS 架构师中心:云端创新的引领者

探索 AWS 架构师中心,获取经实战验证的最佳实践与架构指南,助您高效构建安全、可靠的云上应用