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

如何将Azure Blob存储中持续追加的JSON文件导入Logstash?

实现Azure Blob存储追加式JSON文件的Logstash持续摄入方案

针对你需要将Azure Blob存储中持续追加的JSON文件以类似本地tail的方式导入Logstash的需求,目前常用且受支持的方案主要有以下三种:

1. 基于Azure Event Hubs的实时触发方案

这是最贴近本地tail实时性的方案,通过Azure原生事件机制追踪Blob的追加更新,将增量数据推送到Event Hubs,再由Logstash消费处理。

步骤与配置:

  • Azure端配置:在目标Blob存储容器中创建事件订阅,将BlobCreated和BlobUpdated事件(针对追加Blob)路由到Azure Event Hubs实例。
  • Logstash配置:使用官方维护的azure_event_hubs输入插件,配置示例如下:
input {
  azure_event_hubs {
    connection_string => "Endpoint=sb://<你的Event Hubs命名空间>.servicebus.windows.net/;SharedAccessKeyName=<访问密钥名称>;SharedAccessKey=<访问密钥>;EntityPath=<Event Hubs名称>"
    consumer_group => "$Default"
    codec => "json"
    type => "azure-blob-json"
  }
}

filter {
  # 可按需添加过滤逻辑,比如提取Blob元数据、清洗JSON字段
}

output {
  elasticsearch {
    hosts => ["<你的ES地址>"]
    index => "azure-blob-logs-%{+YYYY.MM.dd}"
  }
}
  • 关键要点:确保Blob存储的事件订阅仅针对追加Blob类型的更新,避免重复触发;Event Hubs的消费组可确保Logstash重启后从断点继续消费,无需重新处理全量数据。

2. 基于http_poller插件的轮询方案

如果不想依赖额外的Azure中间服务,可以使用Logstash内置的http_poller插件,通过Azure Blob Storage REST API定期轮询并获取Blob的增量内容。

配置示例:

input {
  http_poller {
    urls => {
      azure_blob => {
        url => "https://<你的存储账户名称>.blob.core.windows.net/<容器名称>/<目标Blob路径>?restype=blob&comp=appendblocklist"
        headers => {
          "Authorization" => "SharedKey <你的存储账户名称>:<生成的签名>"
          "x-ms-version" => "2021-06-08"
        }
      }
    }
    request_timeout => 60
    interval => 30 # 每30秒轮询一次,可按需调整
    codec => "json"
    type => "azure-blob-json"
    metadata_target => "http_poller_metadata"
  }
}

filter {
  ruby {
    code => "
      # 读取本地记录的已处理块ID,追踪增量
      last_processed_block = File.read('/path/to/sincedb/azure_blob_since').strip rescue ''
      new_blocks = event.get('blocks').select { |b| b['name'] > last_processed_block }
      # 调用Azure API获取新增块内容并合并为JSON数组(需补充具体逻辑)
      File.write('/path/to/sincedb/azure_blob_since', new_blocks.last['name']) if new_blocks.any?
    "
  }
}

output {
  elasticsearch {
    hosts => ["<你的ES地址>"]
    index => "azure-blob-logs-%{+YYYY.MM.dd}"
  }
}
  • 关键要点:利用Azure Blob的追加块列表API追踪增量;通过自定义sincedb文件记录已处理的块ID,避免重复消费;轮询间隔可根据业务延迟需求调整。

3. 基于Azure Data Factory (ADF) 的集成方案

如果需要更复杂的数据转换或多源集成,可使用ADF作为中间层,由ADF监控Blob的新增内容,将增量数据推送到Logstash可消费的端点(如HTTP、Event Hubs)。

核心流程:

  1. 在ADF中创建Blob触发器,监听目标容器的追加Blob更新。
  2. 创建数据管道,将Blob中的增量JSON数据读取后,通过Web活动推送到Logstash的HTTP输入端点。
  3. Logstash配置http输入插件接收数据,示例:
input {
  http {
    port => 8080
    codec => "json"
    type => "azure-blob-json"
  }
}
  • 关键要点:ADF可自动处理增量追踪(通过水印机制),无需手动维护sincedb;适合需要在摄入前进行数据清洗、格式转换的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 02:36:02