Flink+Hologres亿级用户实时UV精确去重最佳实践

简介: Flink+Hologres亿级用户实时UV精确去重最佳实践

UV、PV计算,因为业务需求不同,通常会分为两种场景:

  • 离线计算场景:以T+1为主,计算历史数据
  • 实时计算场景:实时计算日常新增的数据,对用户标签去重

针对离线计算场景,Hologres基于RoaringBitmap,提供超高基数的UV计算,只需进行一次最细粒度的预聚合计算,也只生成一份最细粒度的预聚合结果表,就能达到亚秒级查询。具体详情可以参见往期文章>>Hologres如何支持超高基数UV计算(基于RoaringBitmap实现)

对于实时计算场景,可以使用Flink+Hologres方式,并基于RoaringBitmap,实时对用户标签去重。这样的方式,可以较细粒度的实时得到用户UV、PV数据,同时便于根据需求调整最小统计窗口(如最近5分钟的UV),实现类似实时监控的效果,更好的在大屏等BI展示。相较于以天、周、月等为单位的去重,更适合在活动日期进行更细粒度的统计,并且通过简单的聚合,也可以得到较大时间单位的统计结果。

主体思想

  1. Flink将流式数据转化为表与维表进行JOIN操作,再转化为流式数据。此举可以利用Hologres维表的insertIfNotExists特性结合自增字段实现高效的uid映射。
  2. Flink把关联的结果数据按照时间窗口进行处理,根据查询维度使用RoaringBitmap进行聚合,并将查询维度以及聚合的uid存放在聚合结果表,其中聚合出的uid结果放入Hologres的RoaringBitmap类型的字段中。
  3. 查询时,与离线方式相似,直接按照查询条件查询聚合结果表,并对其中关键的RoaringBitmap字段做or运算后并统计基数,即可得出对应用户数。
  4. 处理流程如下图所示

0.jpeg


方案最佳实践

1.创建相关基础表

1)创建表uid_mapping为uid映射表,用于映射uid到32位int类型。

  • RoaringBitmap类型要求用户ID必须是32位int类型且越稠密越好(即用户ID最好连续)。常见的业务系统或者埋点中的用户ID很多是字符串类型或Long类型,因此需要使用uid_mapping类型构建一张映射表。映射表利用Hologres的SERIAL类型(自增的32位int)来实现用户映射的自动管理和稳定映射。
  • 由于是实时数据, 设置该表为行存表,以提高Flink维表实时JOIN的QPS。
BEGIN;CREATETABLE public.uid_mapping(uid textNOTNULL,uid_int32 serial,PRIMARY KEY (uid));--将uid设为clustering_key和distribution_key便于快速查找其对应的int32值CALL set_table_property('public.uid_mapping','clustering_key','uid');CALL set_table_property('public.uid_mapping','distribution_key','uid');CALL set_table_property('public.uid_mapping','orientation','row');COMMIT;


2)创建表dws_app为基础聚合表,用于存放在基础维度上聚合后的结果。

  • 使用RoaringBitmap前需要创建RoaringBitmap extention,同时也需要Hologres实例为0.10版本
CREATE EXTENSION IF NOT EXISTS roaringbitmap;
  • 为了更好性能,建议根据基础聚合表数据量合理的设置Shard数,但建议基础聚合表的Shard数设置不超过计算资源的Core数。推荐使用以下方式通过Table Group来设置Shard数
--新建shard数为16的Table Group,--因为测试数据量百万级,其中后端计算资源为100core,设置shard数为16BEGIN;CREATETABLE tg16 (a int);--Table Group哨兵表call set_table_property('tg16','shard_count','16');COMMIT;
  • 相比离线结果表,此结果表增加了时间戳字段,用于实现以Flink窗口周期为单位的统计。结果表DDL如下:
BEGIN;createtable dws_app(  country text,  prov text,  city text,  ymd textNOTNULL,--日期字段  timetz TIMESTAMPTZ,--统计时间戳,可以实现以Flink窗口周期为单位的统计  uid32_bitmap roaringbitmap,-- 使用roaringbitmap记录uv  primary key(country, prov, city, ymd, timetz)--查询维度和时间作为主键,防止重复插入数据);CALL set_table_property('public.dws_app','orientation','column');--日期字段设为clustering_key和event_time_column,便于过滤CALL set_table_property('public.dws_app','clustering_key','ymd');CALL set_table_property('public.dws_app','event_time_column','ymd');--等价于将表放在shard数为16的table groupcall set_table_property('public.dws_app','colocate_with','tg16');--group by字段设为distribution_keyCALL set_table_property('public.dws_app','distribution_key','country,prov,city');COMMIT;

2.Flink实时读取数据并更新dws_app基础聚合表

完整示例源码请见alibabacloud-hologres-connectors examples

1)Flink 流式读取数据源(DataStream),并转化为源表(Table)

//此处使用csv文件作为数据源,也可以是kafka等DataStreamSourceodsStream=env.createInput(csvInput, typeInfo);
// 与维表join需要添加proctime字段,详见https://help.aliyun.com/document_detail/62506.htmlTableodsTable=tableEnv.fromDataStream(
odsStream,
$("uid"),
$("country"),
$("prov"),
$("city"),
$("ymd"),
$("proctime").proctime());
// 注册到catalog环境tableEnv.createTemporaryView("odsTable", odsTable);

2)将源表与Hologres维表(uid_mapping)进行关联

其中维表使用insertIfNotExists参数,即查询不到数据时自行插入,uid_int32字段便可以利用Hologres的serial类型自增创建。

// 创建Hologres维表,其中nsertIfNotExists表示查询不到则自行插入StringcreateUidMappingTable=String.format(
"create table uid_mapping_dim("+"  uid string,"+"  uid_int32 INT"+") with ("+"  'connector'='hologres',"+"  'dbname' = '%s',"//Hologres DB名+"  'tablename' = '%s',"//Hologres 表名+"  'username' = '%s',"//当前账号access id+"  'password' = '%s',"//当前账号access key+"  'endpoint' = '%s',"//Hologres endpoint+"  'insertifnotexists'='true'"+")",
database, dimTableName, username, password, endpoint);
tableEnv.executeSql(createUidMappingTable);
// 源表与维表joinStringodsJoinDim="SELECT ods.country, ods.prov, ods.city, ods.ymd, dim.uid_int32"+"  FROM odsTable AS ods JOIN uid_mapping_dim FOR SYSTEM_TIME AS OF ods.proctime AS dim"+"  ON ods.uid = dim.uid";
TablejoinRes=tableEnv.sqlQuery(odsJoinDim);


3)将关联结果转化为DataStream,通过Flink时间窗口处理,结合RoaringBitmap进行聚合

DataStream<Tuple6<String, String, String, String, Timestamp, byte[]>>processedSource=source// 筛选需要统计的维度(country, prov, city, ymd)    .keyBy(0, 1, 2, 3)
// 滚动时间窗口;此处由于使用读取csv模拟输入流,采用ProcessingTime,实际使用中可使用EventTime    .window(TumblingProcessingTimeWindows.of(Time.minutes(5)))
// 触发器,可以在窗口未结束时获取聚合结果    .trigger(ContinuousProcessingTimeTrigger.of(Time.minutes(1)))
    .aggregate(
// 聚合函数,根据key By筛选的维度,进行聚合newAggregateFunction<Tuple5<String, String, String, String, Integer>,
RoaringBitmap,
RoaringBitmap>() {
@OverridepublicRoaringBitmapcreateAccumulator() {
returnnewRoaringBitmap();
            }
@OverridepublicRoaringBitmapadd(
Tuple5<String, String, String, String, Integer>in,
RoaringBitmapacc) {
// 将32位的uid添加到RoaringBitmap进行去重acc.add(in.f4);
returnacc;
            }
@OverridepublicRoaringBitmapgetResult(RoaringBitmapacc) {
returnacc;
            }
@OverridepublicRoaringBitmapmerge(
RoaringBitmapacc1, RoaringBitmapacc2) {
returnRoaringBitmap.or(acc1, acc2);
            }
     },
//窗口函数,输出聚合结果newWindowFunction<RoaringBitmap,
Tuple6<String, String, String, String, Timestamp, byte[]>,
Tuple,
TimeWindow>() {
@Overridepublicvoidapply(
Tuplekeys,
TimeWindowtimeWindow,
Iterable<RoaringBitmap>iterable,
Collector<Tuple6<String, String, String, String, Timestamp, byte[]>>out)
throwsException {
RoaringBitmapresult=iterable.iterator().next();
// 优化RoaringBitmapresult.runOptimize();
// 将RoaringBitmap转化为字节数组以存入Holo中byte[] byteArray=newbyte[result.serializedSizeInBytes()];
result.serialize(ByteBuffer.wrap(byteArray));
// 其中 Tuple6.f4(Timestamp) 字段表示以窗口长度为周期进行统计,以秒为单位out.collect(
newTuple6<>(
keys.getField(0),
keys.getField(1),
keys.getField(2),
keys.getField(3),
newTimestamp(
timeWindow.getEnd() /1000*1000),
byteArray));
        }
    });


4)写入结果表

需要注意的是,Hologres中RoaringBitmap类型在Flink中对应Byte数组类型

// 计算结果转换为表TableresTable=tableEnv.fromDataStream(
processedSource,
$("country"),
$("prov"),
$("city"),
$("ymd"),
$("timest"),
$("uid32_bitmap"));
// 创建Hologres结果表, 其中Hologres的RoaringBitmap类型通过Byte数组存入StringcreateHologresTable=String.format(
"create table sink("+"  country string,"+"  prov string,"+"  city string,"+"  ymd string,"+"  timetz timestamp,"+"  uid32_bitmap BYTES"+") with ("+"  'connector'='hologres',"+"  'dbname' = '%s',"+"  'tablename' = '%s',"+"  'username' = '%s',"+"  'password' = '%s',"+"  'endpoint' = '%s',"+"  'connectionSize' = '%s',"+"  'mutatetype' = 'insertOrReplace'"+")",
database, dwsTableName, username, password, endpoint, connectionSize);
tableEnv.executeSql(createHologresTable);
// 写入计算结果到dws表tableEnv.executeSql("insert into sink select * from "+resTable);

3.数据查询

查询时,从基础聚合表(dws_app)中按照查询维度做聚合计算,查询bitmap基数,得出group by条件下的用户数

  • 查询某天内各个城市的uv
--运行下面RB_AGG运算查询,可执行参数先关闭三阶段聚合开关(默认关闭),性能更好sethg_experimental_enable_force_three_stage_agg=offSELECTcountry        ,prov        ,city        ,RB_CARDINALITY(RB_OR_AGG(uid32_bitmap)) ASuvFROMdws_appWHEREymd='20210329'GROUPBYcountry         ,prov         ,city;


  • 查询某段时间内各个省份的uv
--运行下面RB_AGG运算查询,可执行参数先关闭三阶段聚合开关(默认关闭),性能更好sethg_experimental_enable_force_three_stage_agg=offSELECTcountry        ,prov        ,RB_CARDINALITY(RB_OR_AGG(uid32_bitmap)) ASuvFROMdws_appWHEREtime>'2021-04-19 18:00:00+08'andtime<'2021-04-19 19:00:00+08'GROUPBYcountry         ,prov;
相关实践学习
基于Hologres+PAI+计算巢,5分钟搭建企业级AI问答知识库
本场景采用阿里云人工智能平台PAI、Hologres向量计算和计算巢,搭建企业级AI问答知识库。通过本教程的操作,5分钟即可拉起大模型(PAI)、向量计算(Hologres)与WebUI资源,可直接进行对话问答。
相关文章
|
22小时前
|
SQL 运维 关系型数据库
Flink+Hologres搭建实时数仓
该方案利用Flink和Hologres构建实时数仓,解决传统数仓中间层查询困难、数据不可复用和架构冗余的问题。Flink负责数据源接入和加工,将数据写入Hologres的ODS、DWD和DWS层。Hologres支持高效更新和查询,各层数据可直接服务,简化架构,提高效率。方案具备高性能(Flink与Hologres深度集成,支持实时写入查询)、高可用(主从实例确保服务稳定)和低运维(全链路Flink SQL,减少运维成本)优势。适用于实时报表、推荐系统和业务监控等场景。
18 4
|
1天前
|
Oracle 关系型数据库 MySQL
实时计算 Flink版操作报错合集之用CTAS从mysql同步数据到hologres,改了字段长度,报错提示需要全部重新同步如何解决
在使用实时计算Flink版过程中,可能会遇到各种错误,了解这些错误的原因及解决方法对于高效排错至关重要。针对具体问题,查看Flink的日志是关键,它们通常会提供更详细的错误信息和堆栈跟踪,有助于定位问题。此外,Flink社区文档和官方论坛也是寻求帮助的好去处。以下是一些常见的操作报错及其可能的原因与解决策略。
42 8
|
2天前
|
安全 Java 数据处理
实时计算 Flink版操作报错合集之hologres里报错:找不到字段如何解决
在使用实时计算Flink版过程中,可能会遇到各种错误,了解这些错误的原因及解决方法对于高效排错至关重要。针对具体问题,查看Flink的日志是关键,它们通常会提供更详细的错误信息和堆栈跟踪,有助于定位问题。此外,Flink社区文档和官方论坛也是寻求帮助的好去处。以下是一些常见的操作报错及其可能的原因与解决策略。
14 4
|
5天前
|
SQL 运维 Cloud Native
基于OceanBase+Flink CDC,云粒智慧实时数仓演进之路
本文讲述了其数据中台在传统数仓技术框架下做的一系列努力后,跨进 FlinkCDC 结合 OceanBase 的实时数仓演进过程。
226 2
 基于OceanBase+Flink CDC,云粒智慧实时数仓演进之路
|
5天前
|
SQL 存储 JSON
Flink+Paimon+Hologres 构建实时湖仓数据分析
本文整理自阿里云高级专家喻良,在 Flink Forward Asia 2023 主会场的分享。
|
5天前
|
SQL 存储 JSON
Flink+Paimon+Hologres 构建实时湖仓数据分析
本文整理自阿里云高级专家喻良,在 Flink Forward Asia 2023 主会场的分享。
71646 4
Flink+Paimon+Hologres 构建实时湖仓数据分析
|
5天前
|
存储 消息中间件 监控
基于 Hologres+Flink 的曹操出行实时数仓建设
本文主要介绍曹操出行实时计算负责人林震,基于 Hologres+Flink 的曹操出行实时数仓建设的解决方案分享。
109429 1
基于 Hologres+Flink 的曹操出行实时数仓建设
|
5天前
|
SQL 关系型数据库 MySQL
使用CTAS 把mysql 表同步数据 到hologres ,Flink有什么参数可以使hologres 的字段都小写吗?
使用CTAS 把mysql 表同步数据 到hologres ,Flink有什么参数可以使hologres 的字段都小写吗?
319 0
|
5天前
|
SQL 消息中间件 Kafka
flink问题之做实时数仓sql保证分topic区有序如何解决
Apache Flink是由Apache软件基金会开发的开源流处理框架,其核心是用Java和Scala编写的分布式流数据流引擎。本合集提供有关Apache Flink相关技术、使用技巧和最佳实践的资源。
717 3
|
5天前
|
存储 运维 监控
飞书深诺基于Flink+Hudi+Hologres的实时数据湖建设实践
通过对各个业务线实时需求的调研了解到,当前实时数据处理场景是各个业务线基于Java服务独自处理的。各个业务线实时能力不能复用且存在计算资源的扩展性问题,而且实时处理的时效已不能满足业务需求。鉴于当前大数据团队数据架构主要解决离线场景,无法承接更多实时业务,因此我们需要重新设计整合,从架构合理性,复用性以及开发运维成本出发,建设一套通用的大数据实时数仓链路。本次实时数仓建设将以游戏运营业务为典型场景进行方案设计,综合业务时效性、资源成本和数仓开发运维成本等考虑,我们最终决定基于Flink + Hudi + Hologres来构建阿里云云原生实时湖仓,并在此文中探讨实时数据架构的具体落地实践。
飞书深诺基于Flink+Hudi+Hologres的实时数据湖建设实践

热门文章

最新文章

相关产品

  • 实时数仓 Hologres