实习手记:日志接入 ES 方案(二)——基于 AutoMQ + Logstash 的跨网隔离日志安全传输实战
背景碎碎念:在上篇博客里,我们搞定了本地 IDC 内网环境下的极简高可用 Logstash 架构。原以为日志接入这块可以暂时告一段落,结果导师在每周架构会上抛出了一个更现实也更棘手的问题:
“我们云上(AWS/阿里云/腾讯云 VPC)部署了一套核心业务,本地 IDC 也有自建的数据中心,两边出于安全合规要求,网络是严格物理/逻辑隔离的。云上的日志怎么在不暴露公网风险、不破坏网络隔离策略的前提下,安全、稳妥地传回本地 IDC 的 Elasticsearch?”
听到这个需求,我意识到上一篇“局域网直连”的逻辑完全失效了。在对比了多种方案后,我们最终敲定了以 AutoMQ(云原生 Kafka 替代方案,存算分离架构) 为跨网缓冲缓冲垫,结合 Logstash 和 Elasticsearch 的企业级日志传输架构。今天就来复盘一下这场跨越“云-IDC 隔离墙”的日志传输实战!
01. 痛点:为什么选择 AutoMQ 配合 Logstash?
在真正的企业级混合云/异构网络架构中,将云上日志传输到本地 IDC 会遇到极度苛刻的网络与安全限制:
02. 基于 AutoMQ 的跨网日志传输整体架构

核心传输链路:
03. 配置与代码实战:AutoMQ 生产级接入与 Logstash 消费
在 AutoMQ 中,它完全兼容 Kafka 2.0+ 的协议,因此在 Logstash 端我们可以直接使用标准的 kafka input 插件进行无缝对接。
步骤 1:AutoMQ 层面创建 Topic 与安全认证 (SASL/PLAIN + TLS)
为了防止公网或跨网传输泄露数据,我们必须在 AutoMQ 上开启 SASL 身份认证与 TLS 加密,并为 Logstash 分配专属的消费权限:
# 在 AutoMQ 控制台或使用 kafka-topics.sh 创建专属日志 Topic
kafka-topics.sh –bootstrap-server automq-cluster.company.com:9093 \\
–command-config client.properties \\
–create –topic cloud-security-logs \\
–partitions 6 –replication-factor 3
步骤 2:本地 IDC 侧 Logstash 消费者配置
本地 IDC 的 Logstash 通过标准 Kafka 协议订阅云上 AutoMQ 的 Topic,配置如下:
# automq_to_es_pipeline.conf (位于本地 IDC)
input {
kafka {
# 指向云上 AutoMQ 的接入点(可走专线 IP 或经过 mTLS 代理暴露的域名)
bootstrap_servers => "automq-broker.company.com:9093"
# 消费组配置,方便后期多台 Logstash 实现横向负载均衡消费
group_id => "idc_logstash_consumer_group"
topics => ["cloud-security-logs"]
# 消息格式解析(预先假设云上 Filebeat 写入的是 JSON)
codec => "json"
# 消费拉取性能调优:提升单次 Batch 的拉取条数与字节数
max_poll_records => "2000"
fetch_max_bytes => "10485760" # 10MB
consumer_threads => 4
# 安全认证配置(SASL_SSL 保证传输加密与身份鉴权)
security_protocol => "SASL_SSL"
sasl_mechanism => "PLAIN"
sasl_jaas_config => 'org.apache.kafka.common.security.plain.PlainLoginModule required username="logstash_reader" password="YourSecurePassword123!";'
# SSL 根证书配置,用于校验 AutoMQ 服务端证书合法性
ssl_truststore_location => "/etc/logstash/certs/automq_truststore.jks"
ssl_truststore_password => "TruststorePassword123"
}
}
filter {
# 1. 补充标识,标注该日志来源于云上 AutoMQ 管道
mutate {
add_field => { "[@metadata][source_pipeline]" => "automq_cloud" }
}
# 2. 时间戳修正与清洗
if [log_timestamp] {
date {
match => [ "log_timestamp", "yyyy-MM-dd HH:mm:ss", "ISO8601" ]
target => "@timestamp"
timezone => "Asia/Shanghai"
}
}
# 3. 敏感信息掩码(如掩码 Token/密钥)
mutate {
gsub => [ "message", "(access_key|secret_key)=[^&]+", "\\1=******" ]
}
}
output {
# 批量写入本地 IDC 的 Elasticsearch 集群
elasticsearch {
hosts => ["http://192.168.10.11:9200", "http://192.168.10.12:9200"]
index => "cloud-sec-log-%{+YYYY.MM.dd}"
user => "logstash_writer"
password => "InternalESPass123!"
# 开启并发 worker 提升 ES 写入吞吐量
workers => 4
}
}
04. 网络断连与突发峰值下的“削峰填谷”机制
为什么这套架构比直接使用 Logstash 直连更加稳健?在网络断连或突发流量洪峰时,AutoMQ 是如何发挥“蓄水池”作用的:

生产级防护优势:
传统 Kafka 积压数据受限于昂贵的 EBS 云盘容量,积压久了容易把磁盘撑爆。而 AutoMQ 的数据大部分直接写入云厂商的 S3/OSS 对象存储,具备低成本、高可靠、无上限的存储能力。即便跨网网络中断整整两天,云上日志也能安全地保存在 AutoMQ 中,绝不丢失。
云上即便产生 10 万 EPS(每秒日志条数)的峰值,本地 IDC 的 Logstash 依然可以根据 ES 集群的写入能力,稳定保持在 1 万 EPS 的速率“匀速拉取”。AutoMQ 在中间完美承担了“蓄水池”与“缓冲垫”的角色。
Logstash 只有在成功把数据 Batch 写入 ES 后,才会向 AutoMQ 提交当前消费的 Offset。如果在传输或解析过程中 Logstash 进程崩溃,重启后会从上一次成功提交的 Offset 继续消费,保证数据零丢失。
05. 实习小结与后续预告
通过引入 AutoMQ →\\rightarrow→ Logstash →\\rightarrow→ Elasticsearch 这套架构,我们成功解决了一系列硬核难题:
- 架构解耦:云上采集端只管往 AutoMQ 吐数据,本地 IDC 只管按需消费,两边完全解耦。
- 抗暴击与零丢失:借助 AutoMQ 的云原生存算分离特性,轻松应对公网断连与突发流量洪峰。
- 安全可控:采用 SASL_SSL 强认证与加密传输,满足企业级的网络安全合规要求。
现在,无论是 IDC 内网的日志,还是云上 VPC 里的隔离日志,都能源源不断、安全无虞地汇聚到我们本地的 Elasticsearch 中,为后面的 AI SOC 平台提供坚实的数据支撑了!


