flink-reduce

栏目: 编程工具 · 发布时间: 7年前

阅读更多

一.背景

有时候我们需要过滤数据,有些中间数据是不需要的,比如场景:

binlog 数据更新的时候,我们仅仅需要最新数据。会根据ID 分组,然后取version 最大的一条,存储

二.简单实例

@Data
@ToString
public class Order {
    // 主键id
    private Integer id;
    // 版本
    private Integer version;
    private Timestamp mdTime;

    public Order(int id, Integer version) {
        this.id = id;
        this.version = version;
        this.mdTime = new Timestamp(System.currentTimeMillis());
    }

    public Order() {
    }
}
public class OrderSource implements SourceFunction<Order> {
    Random random = new Random();

    @Override
    public void run(SourceContext<Order> ctx) throws Exception {
        while (true) {
            TimeUnit.MILLISECONDS.sleep(100);
            // 为了区分,我们简单生0~2的id, 和版本0~99
            int id = random.nextInt(3);
            Order o = new Order(id, random.nextInt(100));
            ctx.collect(o);
        }
    }
    @Override
    public void cancel() {

    }
}
public class ReduceApp {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
        DataStream<Order> userInfoDataStream = env.addSource(new OrderSource());

        DataStream<Order> timedData = userInfoDataStream.assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Order>() {
            @Override
            public long extractAscendingTimestamp(Order element) {
                return element.getMdTime().getTime();
            }
        });

        SingleOutputStreamOperator<Order> reduce = timedData
                .keyBy("id")
                .timeWindow(Time.seconds(10), Time.seconds(5))
                .reduce((ReduceFunction<Order>) (v1, v2) -> v1.getVersion() >= v2.getVersion() ? v1 : v2);

        reduce.print();
        env.execute("test");
    }
}

结果:

Order(id=2, version=97, mdTime=2019-03-11 17:39:34.052)

Order(id=0, version=99, mdTime=2019-03-11 17:39:32.913)

Order(id=1, version=96, mdTime=2019-03-11 17:39:34.155)

Order(id=2, version=97, mdTime=2019-03-11 17:39:34.052)

Order(id=1, version=96, mdTime=2019-03-11 17:39:34.155)

Order(id=0, version=99, mdTime=2019-03-11 17:39:32.913)

这个会对同一个窗口做过滤,比如同步到另一个mysql,hdfs,就能减少数据量

0顶

0踩

分享到:

flink-watermark

评论


以上所述就是小编给大家介绍的《flink-reduce》,希望对大家有所帮助,如果大家有任何疑问请给我留言,小编会及时回复大家的。在此也非常感谢大家对 码农网 的支持!

查看所有标签

猜你喜欢:

本站部分资源来源于网络,本站转载出于传递更多信息之目的,版权归原作者或者来源机构所有,如转载稿涉及版权问题,请联系我们

编程之法

编程之法

July / 人民邮电出版社 / 2015-9-1 / 49.00元

本书涉及面试、算法、机器学习三个主题。书中的每道编程题目都给出了多种思路、多种解法,不断优化、逐层递进。本书第1章至第6章分别阐述字符串、数组、树、查找、动态规划、海量数据处理等相关的编程面试题和算法,第7章介绍机器学习的两个算法—K近邻和SVM。此外,每一章都有“举一反三”和“习题”,以便读者及时运用所学的方法解决相似的问题,且在附录中收录了语言、链表、概率等其他题型。书中的每一道题都是面试的高......一起来看看 《编程之法》 这本书的介绍吧!

图片转BASE64编码
图片转BASE64编码

在线图片转Base64编码工具

正则表达式在线测试
正则表达式在线测试

正则表达式在线测试

RGB HSV 转换
RGB HSV 转换

RGB HSV 互转工具