Advertisement

【flink】udf实现维表关联

阅读量:

需求

复制代码
    	用udf实现flink维表关联。
    	维表数据库用的是Cassandra,具有快速高并发写入的特点,把明细数据保存在数据库中。主流是接受kafka消息(本文中用的datagen),根据自己的过滤条件(where xx=xxx)在udf中查询出一批数据,做聚合后把结果返回。
    	本文中udf的功能为,查询近5分钟的count 和sum。调用udf后主流就有了这两个属性。

代码

复制代码
    /** * 统计5分钟的 count sum
     */
    public class AggCassandraFunction extends ScalarFunction {
    private ClusterBuilder builder;
    private CassandraInputFormat<Tuple3<String, Integer, Integer>> format;
    
    @Override
    public void open(FunctionContext context) throws Exception {
        builder = new ClusterBuilder() {
    
            protected Cluster bui

全部评论 (0)

还没有任何评论哟~