【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)
还没有任何评论哟~
