【问题标题】:Event from Netcat source doesn't go through Kafka channel来自 Netcat 源的事件不通过 Kafka 通道
【发布时间】:2016-03-25 16:19:05
【问题描述】:

我使用 Flume 代理通过 Flume 代理收集外部数据。外部数据批处理几乎是每 10 秒 1MB。我将 Flume 代理配置如下。

# Flume agent configuration as /flume/conf/agent.conf
agent.sources = netcat-source
agent.channels = kafka-channel
agent.sinks = logger-sink

########################################
#   Netcat Source
########################################

agent.sources.netcat-source.type = netcat
agent.sources.netcat-source.bind = 0.0.0.0
agent.sources.netcat-source.port = 4141
agent.sources.netcat-source.max-line-length = 500000
agent.sources.netcat-source.channels = kafka-channel

########################################
#   Kafka Channel
########################################

agent.channels.kafka-channel.type =  org.apache.flume.channel.kafka.KafkaChannel
agent.channels.kafka-channel.brokerList = 10.212.136.108:9092,10.212.136.108:9092
agent.channels.kafka-channel.zookeeperConnect = 10.212.136.108:2181,10.212.136.108:2181/kafka
agent.channels.kafka-channel.topic = channel
agent.channels.kafka-channel.groupId = fcd-group


########################################
#   Logger Sink
########################################

agent.sinks.logger-sink.type = logger
agent.sinks.logger-sink.channel = kafka-channel

我通过以下方式激活了代理。

flume-ng agent -n agent -c /flume/conf -f /flume/conf/agent.conf 

不幸的是,netcat 源运行良好,但通道或接收器出现问题。从 Ubuntu 的资源监视器中,我可以看到以下性能。 Network performance. Blue curve indicates input while red one indicates output 如果没有其他应用程序与网络 io 一起运行,我相信这个图展示了我的 Flume 代理发生了什么。

当我通过控制台消费者检查主题“频道”中的 Kafka 内容时,我什么也没得到。另外,当我检查 flume.log 时,我只得到了 Flume 的状态输出,没有数据。

我已经使用

验证了传入的数据
nc -lk 4141 >> my_data_check_file

我的频道或接收器出了什么问题?

附:当我使用内存通道、文件通道时,事情变得同样棘手。

【问题讨论】:

    标签: logging apache-kafka flume


    【解决方案1】:

    啊,终于,我自己解决了这个问题!

    关键点是行分隔符'\n'。

    在 Flume 源代码 NetcatSource.java 中,我们有一个像下面这样的复杂行

    private int processEvents(CharBuffer buffer, Writer writer) throws IOException {
      int numProcessed = 0;
    
      boolean foundNewLine = true;
      while (foundNewLine) {
        foundNewLine = false;
    
        int limit = buffer.limit();
        for (int pos = buffer.position(); pos < limit; pos++) {
          if (buffer.get(pos) == '\n') {  
            // parse event body bytes out of CharBuffer
            buffer.limit(pos); // temporary limit
            ByteBuffer bytes = Charsets.UTF_8.encode(buffer);
            buffer.limit(limit); // restore limit
    ... ...
    ... ...
    

    代码强制输入数据以“\n”结尾。否则,通道将不会处理任何事件。我们可以根据需要更改这个字符并将自定义的源放入 $FLUME_HOME/lib

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-05-27
      • 2018-12-15
      • 1970-01-01
      • 1970-01-01
      • 2018-01-26
      • 2016-07-21
      相关资源
      最近更新 更多