Flink中的广播流之BroadcastStream

使用场景:
在处理数据的时候,有些配置是要实时动态改变的,比如说我要过滤一些关键字,这些关键字呢是在MYSQL里随时配置修改的,那我们在高吞吐计算的Function中动态查询配置文件有可能使整个计算阻塞,甚至任务停止。
广播流可以通过查询配置文件,广播到某个 operator 的所有并发实例中,然后与另一条流数据连接进行计算。

实现步骤:
1、定义一个MapStateDescriptor来描述我们要广播的数据的格式

final MapStateDescriptor<String, String> CONFIG_DESCRIPTOR = new MapStateDescriptor<>(
                "wordsConfig",
                BasicTypeInfo.STRING_TYPE_INFO,
                BasicTypeInfo.STRING_TYPE_INFO);

2、需要一个Stream来广播下游的operator
我这里实现了一个只有1个并发度的数据源,定时查配置文件,发动到下游

public class MinuteBroadcastSource extends RichParallelSourceFunction<String> {


    private volatile boolean isRun;
    private volatile int lastUpdateMin = -1;

    private R2mClusterClient redisDao;

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        isRun = true;
    }

    @Override
    public void run(SourceContext<String> ctx) throws Exception {
        while(isRun){
            LocalDateTime date = LocalDateTime.now();
            int min = date.getMinute();
            if(min != lastUpdateMin){
                lastUpdateMin = min;
                Set<String> configs = readConfigs();
                if(configs != null && configs.size() > 0){
                    for(String config : configs){
                        ctx.collect(config);
                    }

                }
            }
            Thread.sleep(1000);
        }
    }

    private Set<String> readConfigs(){
        //这里读取配置信息
        return null;
    }


    @Override
    public void cancel() {
        isRun = false;
    }
}

3、添加数据源并把数据源注册成广播流

 BroadcastStream<String> broadcastStream = env.addSource(new MinuteBroadcastSource()).setParallelism(1).broadcast(CONFIG_DESCRIPTOR);

4、连接广播流和处理数据的流

DataStream<SkuOrder> connectedStream = orderStream.connect(broadcastStream).process(new BroadcastProcessFunction<Order, String, Order>(){
            @Override
            public void processElement(Order order, ReadOnlyContext ctx, Collector<Order> collector) throws Exception {
                HeapBroadcastState<String,String> config = (HeapBroadcastState)ctx.getBroadcastState(CONFIG_DESCRIPTOR);
                Iterator<Map.Entry<String, String>> iterator = config.iterator();
                while (iterator.hasNext()){
                    Map.Entry<String, String> entry =iterator.next();
                    logger.info("all config:" + entry.getKey()  +  "   value:" + entry.getValue());
                }
                collector.collect(order);
            }

            @Override
            public void processBroadcastElement(String s, Context ctx, Collector<SkuOrder> collector) throws Exception {
                logger.info("收到广播:" + s);
                BroadcastState<String,String> state =  ctx.getBroadcastState(CONFIG_DESCRIPTOR);
                ctx.getBroadcastState(CONFIG_DESCRIPTOR).put(s,"1"); 
            }
        });

需要注意到的问题:
1、数据源发送数据时候如果数据是集合,必须使用线程安全的集合类
2、如果上面的MinuteBroadcastSource并行度大于1,那么每一个JOB都会发一条广播,这样的话如果每个JOB一分钟发一次,那么processBroadcastElement会收到 并行度数 * n条消息
3、获取到的BroadcastState是一个map,相同的KEY,put进去会覆盖掉

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 194,088评论 5 459
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 81,715评论 2 371
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 141,361评论 0 319
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 52,099评论 1 263
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 60,987评论 4 355
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 46,063评论 1 272
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 36,486评论 3 381
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 35,175评论 0 253
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 39,440评论 1 290
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 34,518评论 2 309
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 36,305评论 1 326
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 32,190评论 3 312
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 37,550评论 3 298
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 28,880评论 0 17
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 30,152评论 1 250
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 41,451评论 2 341
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 40,637评论 2 335

推荐阅读更多精彩内容