Flume配置Failover Sink Processor

Stella981
• 阅读 742

1 官网内容

Flume配置Failover Sink Processor

2 看一张图一目了然

Flume配置Failover Sink Processor

Flume配置Failover Sink Processor

3 详细配置

  source配置文件

#配置文件:
    a1.sources= r1
    a1.sinks= k1 k2
    a1.channels= c1
    
    #负载平衡
    a1.sinkgroups = g1
    a1.sinkgroups.g1.sinks = k1 k2
    a1.sinkgroups.g1.processor.type = failover
    a1.sinkgroups.g1.processor.priority.k1 = 5
    a1.sinkgroups.g1.processor.priority.k2 = 10
    a1.sinkgroups.g1.processor.maxpenalty = 1000
    
    
    
    #Describe/configure the source
    a1.sources.r1.type= exec
    a1.sources.r1.command= tail -F /tmp/logs/test.log
    
    
    #Describe the sink
    a1.sinks.k1.type= avro
    a1.sinks.k1.hostname= 127.0.0.1
    a1.sinks.k1.port= 50001
    
    a1.sinks.k2.type= avro
    a1.sinks.k2.hostname= 127.0.0.1
    a1.sinks.k2.port= 50002
    
    # Usea channel which buffers events in memory
    a1.channels.c1.type= memory
    a1.channels.c1.capacity= 1000
    a1.channels.c1.transactionCapacity= 100
    
    # set channel
    a1.sinks.k1.channel= c1
    a1.sinks.k2.channel= c1
    a1.sources.r1.channels= c1

  sink1配置文件

# Name the components on this agent
    a2.sources = r1
    a2.sinks = k1
    a2.channels = c1
    
    # Describe/configure the source
    a2.sources.r1.type = avro
    a2.sources.r1.channels = c1
    a2.sources.r1.bind = 127.0.0.1
    a2.sources.r1.port = 50001
    
    # Describe the sink
    a2.sinks.k1.type = logger
    a2.sinks.k1.channel = c1
    
    # Use a channel which buffers events inmemory
    a2.channels.c1.type = memory
    a2.channels.c1.capacity = 1000
    a2.channels.c1.transactionCapacity = 100

  sink2配置

# Name the components on this agent
    a3.sources = r1
    a3.sinks = k1
    a3.channels = c1
    
    # Describe/configure the source
    a3.sources.r1.type = avro
    a3.sources.r1.channels = c1
    a3.sources.r1.bind = 127.0.0.1
    a3.sources.r1.port = 50002
    
    # Describe the sink
    a3.sinks.k1.type = logger
    a3.sinks.k1.channel = c1
    
    # Use a channel which buffers events inmemory
    a3.channels.c1.type = memory
    a3.channels.c1.capacity = 1000
    a3.channels.c1.transactionCapacity = 100

4 启动服务

先启动sink1 sink2 再启动source
            
    flume-ng agent -c conf -f /mnt/software/flume-1.6.0/flume-conf/failOver/sink2.conf -n a3 -Dflume.root.logger=DEBUG,console
    flume-ng agent -c conf -f /mnt/software/flume-1.6.0/flume-conf/failOver/sink1.conf -n a2 -Dflume.root.logger=DEBUG,console
    flume-ng agent -c conf -f /mnt/software/flume-1.6.0/flume-conf/failOver/load_source_case.conf -n a1 -Dflume.root.logger=DEBUG,console

  5 效果测试

启动后第一次走了sink2
  
           : /127.0.0.1:42828
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 68 61 64 6F 6F 70                               hadoop }
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 7A 68 61 6E 67 6A 69 6E                         zhangjin }
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 78 78 78 78                                     xxxx }
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 79 79 79 79                                     yyyy }
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 7A 68 61 6E 67 6A 69 6E                         zhangjin }
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 78 78 78 78                                     xxxx }
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 79 79 79 79                                     yyyy }
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 7A 68 61 6E 67 6A 69 6E                         zhangjin }
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 78 78 78 78                                     xxxx }
    19/02/21 23:45:47 INFO sink.LoggerSink: Event: { headers:{} body: 79 79 79 79                                     yyyy }        
    
挂掉sink2,之后source感知到sink2挂了
    
    Caused by: java.net.ConnectException: Connection refused
        at sun.nio.ch.SocketChannelImpl.checkConnect(Native Method)
        at sun.nio.ch.SocketChannelImpl.finishConnect(SocketChannelImpl.java:717)
        at org.jboss.netty.channel.socket.nio.NioClientBoss.connect(NioClientBoss.java:148)
        at org.jboss.netty.channel.socket.nio.NioClientBoss.processSelectedKeys(NioClientBoss.java:104)
        at org.jboss.netty.channel.socket.nio.NioClientBoss.process(NioClientBoss.java:78)
        at org.jboss.netty.channel.socket.nio.AbstractNioSelector.run(AbstractNioSelector.java:312)
        at org.jboss.netty.channel.socket.nio.NioClientBoss.run(NioClientBoss.java:41)
        ... 3 more    
        
        
数据发往sink1
    
        19/02/21 23:45:41 INFO ipc.NettyServer: [id: 0x77bfe0b5, /127.0.0.1:47142 => /127.0.0.1:50001] BOUND: /127.0.0.1:50001
    19/02/21 23:45:41 INFO ipc.NettyServer: [id: 0x77bfe0b5, /127.0.0.1:47142 => /127.0.0.1:50001] CONNECTED: /127.0.0.1:47142
    19/02/21 23:47:14 INFO sink.LoggerSink: Event: { headers:{} body: 7A 68 61 6E 67 6A 69 6E                         zhangjin }
    19/02/21 23:47:14 INFO sink.LoggerSink: Event: { headers:{} body: 78 78 78 78                                     xxxx }
    19/02/21 23:47:14 INFO sink.LoggerSink: Event: { headers:{} body: 79 79 79 79                                     yyyy }

  6 总结,从效果来看sink2挂了之后,数据发往sink1,实现了失败迁移的功能。

点赞
收藏
评论区
推荐文章
blmius blmius
3年前
MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1
文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s
皕杰报表之UUID
​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为
待兔 待兔
6个月前
手写Java HashMap源码
HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22
Jacquelyn38 Jacquelyn38
3年前
2020年前端实用代码段,为你的工作保驾护航
有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )
Wesley13 Wesley13
3年前
mysql设置时区
mysql设置时区mysql\_query("SETtime\_zone'8:00'")ordie('时区设置失败,请联系管理员!');中国在东8区所以加8方法二:selectcount(user\_id)asdevice,CONVERT\_TZ(FROM\_UNIXTIME(reg\_time),'08:00','0
Wesley13 Wesley13
3年前
00:Java简单了解
浅谈Java之概述Java是SUN(StanfordUniversityNetwork),斯坦福大学网络公司)1995年推出的一门高级编程语言。Java是一种面向Internet的编程语言。随着Java技术在web方面的不断成熟,已经成为Web应用程序的首选开发语言。Java是简单易学,完全面向对象,安全可靠,与平台无关的编程语言。
Stella981 Stella981
3年前
Django中Admin中的一些参数配置
设置在列表中显示的字段,id为django模型默认的主键list_display('id','name','sex','profession','email','qq','phone','status','create_time')设置在列表可编辑字段list_editable
Wesley13 Wesley13
3年前
MySQL部分从库上面因为大量的临时表tmp_table造成慢查询
背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_
Python进阶者 Python进阶者
1年前
Excel中这日期老是出来00:00:00,怎么用Pandas把这个去除
大家好,我是皮皮。一、前言前几天在Python白银交流群【上海新年人】问了一个Pandas数据筛选的问题。问题如下:这日期老是出来00:00:00,怎么把这个去除。二、实现过程后来【论草莓如何成为冻干莓】给了一个思路和代码如下:pd.toexcel之前把这
Oracle 分组与拼接字符串同时使用
SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(