基于 Confluent + Flink 的实时数据分析最佳实践

本文涉及的产品
实时计算 Flink 版,5000CU*H 3个月
简介: 在实际业务使用中,需要经常实时做一些数据分析,包括实时PV和UV展示,实时销售数据,实时店铺UV以及实时推荐系统等,基于此类需求,Confluent+实时计算Flink版是一个高效的方案。

业务背景

在实际业务使用中,需要经常实时做一些数据分析,包括实时PV和UV展示,实时销售数据,实时店铺UV以及实时推荐系统等,基于此类需求,Confluent+实时计算Flink版是一个高效的方案。


Confluent是基于Apache Kafka提供的企业级全托管流数据服务,由 Apache Kafka 的原始创建者构建,通过企业级功能扩展了 Kafka 的优势,同时消除了 Kafka管理或监控的负担。


实时计算Flink版是阿里云基于 Apache Flink 构建的企业级实时大数据计算商业产品。实时计算 Flink 由 Apache Flink 创始团队官方出品,拥有全球统一商业化品牌,提供全系列产品矩阵,完全兼容开源 Flink API,并充分基于强大的阿里云平台提供云原生的 Flink 商业增值能力。


一、准备工作-创建Confluent集群和实时计算Flink版集群

  1. 登录Confluent管理控制台,创建Confluent集群,创建步骤参考 Confluent集群开通


  1. 登录实时计算Flink版管理控制台,创建vvp集群。请注意,创建vvp集群选择的vpc跟confluent集群的region和vpc使用同一个,这样可以在vvp内部访问confluent的内部域名。


二、最佳实践-实时统计玩家充值金额-Confluent+实时计算Flink+Hologres

2.1 新建Confluent消息队列

  1. 在confluent集群列表页,登录control center

  1. 在左侧选中Topics,点击Add a topic按钮,创建一个名为confluent-vvp-test的topic,将partition设置为3

2.2 配置结果表 Hologres

  1. 进入Hologres控制台,点击Hologres实例,在DB管理中新增数据库`mydb`

  1. 登录Hologres数据库,新建SQL

  1. Hologres中创建结果表 SQL语句
--用户累计消费结果表
 CREATE TABLE consume (
    appkey VARCHAR,
    serverid VARCHAR,
    servertime VARCHAR,
    roleid VARCHAR,
    amount FLOAT,
    dt VARCHAR,
    primary key(appkey,dt)
  );


2.3 创建实时计算vvp作业

  1. 首先登录vvp控制台,选择集群所在region,点击控制台,进入开发界面

  1. 点击作业开发Tab,点击新建文件,文件名称:confluent-vvp-hologres,文件类型选择:流作业/SQL

  1. 在输入框写入以下代码:
create TEMPORARY table kafka_game_consume_source(  
  appkey STRING,
  servertime STRING,
  consumenum DOUBLE,
  roleid STRING,
  serverid STRING    
) with (
   'connector' = 'kafka',
   'topic' = 'game_consume_log',
   'properties.bootstrap.servers' = 'kafka.confluent.svc.cluster.local.xxx:9071[xxx可以找开发同学查看]',
   'properties.group.id' = 'gamegroup',
   'format' = 'json',
   'properties.ssl.truststore.location' = '/flink/usrlib/truststore.jks',
   'properties.ssl.truststore.password' = '[your truststore password]',
   'properties.security.protocol'='SASL_SSL',
   'properties.sasl.mechanism'='PLAIN',
   'properties.sasl.jaas.config'='org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="xxx[集群的用户]" password="xxx[相应的密码]";'
);
-- 创建累计消费hologres sink表
CREATE TEMPORARY TABLE consume(
 appkey STRING,
   serverid STRING,
  servertime STRING,
  roleid STRING,
  amount DOUBLE,
  dt STRING,
  PRIMARY KEY (appkey,dt) NOT ENFORCED
  )WITH (
  'connector' = 'hologres',
  'dbname' = 'mydb',
  'endpoint' = 'hgprecn-cn-tl32gkaet006-cn-beijing-vpc.hologres.aliyuncs.com:80',
  'password' = '[your appkey secret]',
  'tablename' = 'consume',
  'username' = '[your app key]',
  'mutateType' = 'insertorreplace'
  );
--{"appkey":"appkey1","servertime":"2020-09-30 14:10:36","consumenum":33.8,"roleid":"roleid1","serverid":"1"}
--{"appkey":"appkey2","servertime":"2020-09-30 14:11:36","consumenum":30.8,"roleid":"roleid2","serverid":"2"}
--{"appkey":"appkey1","servertime":"2020-09-30 14:13:36","consumenum":31.8,"roleid":"roleid1","serverid":"1"}
--{"appkey":"appkey2","servertime":"2020-09-30 14:20:36","consumenum":33.8,"roleid":"roleid2","serverid":"2"}
--{"appkey":"appkey1","servertime":"2020-09-30 14:30:36","consumenum":73.8,"roleid":"roleid1","serverid":"1"}
  -- 计算每个用户累积消费金额
  insert into consume
  SELECT
   appkey,LAST_VALUE(serverid) as serverid,LAST_VALUE(servertime) as servertime,LAST_VALUE(roleid) as roleid,
   sum(consumenum) as amount,
  substring(servertime,1,10) as dt
  FROM kafka_game_consume_source
  GROUP BY appkey,substring(servertime,1,10)
  having sum(consumenum) > 0;
  1. 在高级配置里,增加依赖文件truststore.jks(访问内部域名得添加这个文件,访问公网域名可以不用),访问依赖文件的固定路径前缀都是/flink/usrlib/(这里就是/flink/usrlib/truststore.jks)


  1. 点击上线按钮,完成上线


  1. 在运维作用列表里找到刚上线的作用,点击启动按钮,等待状态更新为running,运行成功。


  1. 在control center的【Topics->Messages】页面,逐条发送测试消息,格式为:
{"appkey":"appkey1","servertime":"2020-09-30 14:10:36","consumenum":33.8,"roleid":"roleid1","serverid":"1"}
{"appkey":"appkey2","servertime":"2020-09-30 14:11:36","consumenum":30.8,"roleid":"roleid2","serverid":"2"}
{"appkey":"appkey1","servertime":"2020-09-30 14:13:36","consumenum":31.8,"roleid":"roleid1","serverid":"1"}
{"appkey":"appkey2","servertime":"2020-09-30 14:20:36","consumenum":33.8,"roleid":"roleid2","serverid":"2"}
{"appkey":"appkey1","servertime":"2020-09-30 14:30:36","consumenum":73.8,"roleid":"roleid1","serverid":"1"}


2.4 查看用户充值金额实时统计效果


三、最佳实践-电商实时PV和UV统计-Confluent+实时计算Flink+RDS

3.1 新建Confluent消息队列

  1. 在confluent集群列表页,登录control center


  1. 在左侧选中Topics,点击Add a topic按钮,创建一个名为pv-uv的topic,将partition设置为3

3.2 创建云数据库RDS结果表

  1. 登录 RDS 管理控制台页面,购买RDS。确保RDS与Flink全托管集群在相同region,相同VPC下

  1. 添加虚拟交换机网段(vswitch IP段)进入RDS白名单,详情参考:设置白名单文档

3.【vswitch IP段】可在 flink的工作空间详情中查询


  1. 在【账号管理】页面创建账号【高权限账号】



  1. 数据库实例下【数据库管理】新建数据库【conflufent_vvp】

  1. 使用系统自带的DMS服务登陆RDS,登录名和密码输入上面创建的高权限账户


  1. 双击【confluent_vvp】数据库,打开SQLConsole,将以下建表语句复制粘贴到 SQLConsole中,创建结果表
CREATE TABLE result_cps_total_summary_pvuv_min(
  summary_date date NOT NULL COMMENT '统计日期',
  summary_min varchar(255) COMMENT '统计分钟',
  pv bigint COMMENT 'pv',
  uv bigint COMMENT 'uv',
  currenttime timestamp COMMENT '当前时间',
  primary key(summary_date,summary_min)
)

3.3 创建实时计算VVP作业

1.【[VVP控制台】新建文件


  1. 在SQL区域输入以下代码:
--数据的订单源表
CREATE TABLE source_ods_fact_log_track_action (
  account_id VARCHAR,
  --用户ID
  client_ip VARCHAR,
  --客户端IP
  client_info VARCHAR,
  --设备机型信息
  platform VARCHAR,
  --系统版本信息
  imei VARCHAR,
  --设备唯一标识
  `version` VARCHAR,
  --版本号
  `action` VARCHAR,
  --页面跳转描述
  gpm VARCHAR,
  --埋点链路
  c_time VARCHAR,
  --请求时间
  target_type VARCHAR,
  --目标类型
  target_id VARCHAR,
  --目标ID
  udata VARCHAR,
  --扩展信息,JSON格式
  session_id VARCHAR,
  --会话ID
  product_id_chain VARCHAR,
  --商品ID串
  cart_product_id_chain VARCHAR,
  --加购商品ID
  tag VARCHAR,
  --特殊标记
  `position` VARCHAR,
  --位置信息
  network VARCHAR,
  --网络使用情况
  p_dt VARCHAR,
  --时间分区天
  p_platform VARCHAR --系统版本信息
) WITH (
   'connector' = 'kafka',
   'topic' = 'game_consume_log',
   'properties.bootstrap.servers' = 'kafka.confluent.svc.cluster.local.c79f69095bc5d4d98b01136fe43e31b93:9071',
   'properties.group.id' = 'gamegroup',
   'format' = 'json',
   'properties.ssl.truststore.location' = '/flink/usrlib/truststore.jks',
   'properties.ssl.truststore.password' = '【your password】',
   'properties.security.protocol'='SASL_SSL',
   'properties.sasl.mechanism'='PLAIN',
   'properties.sasl.jaas.config'='org.apache.flink.kafka.shaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="【your user name】" password="【your password】";'
);
--{"account_id":"id1","client_ip":"172.11.1.1","client_info":"mi10","p_dt":"2021-12-01","c_time":"2021-12-01 19:10:00"}
CREATE TABLE result_cps_total_summary_pvuv_min (
  summary_date date,
  --统计日期
  summary_min varchar,
  --统计分钟
  pv bigint,
  --点击量
  uv bigint,
  --一天内同个访客多次访问仅计算一个UV
  currenttime timestamp,
  --当前时间
  primary key (summary_date, summary_min)
) WITH (
  type = 'rds',
  url = 'url = 'jdbc:mysql://rm-【your rds clusterId】.mysql.rds.aliyuncs.com:3306/confluent_vvp',',
  tableName = 'result_cps_total_summary_pvuv_min',
  userName = 'flink_confluent_vip',
  password = '【your rds password】'
);
CREATE VIEW result_cps_total_summary_pvuv_min_01 AS
select
  cast (p_dt as date) as summary_date --时间分区
  , count (client_ip) as pv --客户端的IP
  , count (distinct client_ip) as uv --客户端去重
  , cast (max (c_time) as TIMESTAMP) as c_time --请求的时间
from
  source_ods_fact_log_track_action
group
  by p_dt;
INSERT
  into result_cps_total_summary_pvuv_min
select
  a.summary_date,
  --时间分区
  cast (DATE_FORMAT (c_time, 'HH:mm') as varchar) as summary_min,
  --取出小时分钟级别的时间
  a.pv,
  a.uv,
  CURRENT_TIMESTAMP as currenttime --当前时间
from
  result_cps_total_summary_pvuv_min_01 AS a;
  1. 点击【上线】之后,在作业运维页面点击启动按钮,直到状态更新为RUNNING状态。


  1. 在control center的【Topics->Messages】页面,逐条发送测试消息,格式为:
{"account_id":"id1","client_ip":"72.11.1.111","client_info":"mi10","p_dt":"2021-12-01","c_time":"2021-12-01 19:11:00"}
{"account_id":"id2","client_ip":"72.11.1.112","client_info":"mi10","p_dt":"2021-12-01","c_time":"2021-12-01 19:12:00"}
{"account_id":"id3","client_ip":"72.11.1.113","client_info":"mi10","p_dt":"2021-12-01","c_time":"2021-12-01 19:13:00"}


3.4 查看PV和UV效果

    可以看出rds数据表的pv和uv会随着发送的消息数据,动态的变化,同时还可以通过【数据可视化】来查看相应的图表信息。


pv图表展示:


uv图表展示:



产品技术咨询

https://survey.aliyun.com/apps/zhiliao/VArMPrZOR  

加入技术交流群

image.png


更多 Flink 相关技术问题,可扫码加入社区钉钉交流群

第一时间获取最新技术文章和社区动态,请关注公众号~

image.png

活动推荐

阿里云基于 Apache Flink 构建的企业级产品-实时计算Flink版现开启活动:

99 元试用 实时计算Flink版(包年包月、10CU)即有机会获得 Flink 独家定制卫衣;另包 3 个月及以上还有 85 折优惠!

了解活动详情:https://www.aliyun.com/product/bigdata/sc

image.png

相关实践学习
基于Hologres轻松玩转一站式实时仓库
本场景介绍如何利用阿里云MaxCompute、实时计算Flink和交互式分析服务Hologres开发离线、实时数据融合分析的数据大屏应用。
Linux入门到精通
本套课程是从入门开始的Linux学习课程,适合初学者阅读。由浅入深案例丰富,通俗易懂。主要涉及基础的系统操作以及工作中常用的各种服务软件的应用、部署和优化。即使是零基础的学员,只要能够坚持把所有章节都学完,也一定会受益匪浅。
相关文章
|
3月前
|
SQL 存储 数据库
Flink + Paimon 数据 CDC 入湖最佳实践
Flink + Paimon 数据 CDC 入湖最佳实践
364 1
|
1月前
|
SQL 关系型数据库 MySQL
Flink CDC + Hudi + Hive + Presto构建实时数据湖最佳实践
Flink CDC + Hudi + Hive + Presto构建实时数据湖最佳实践
144 0
|
1月前
|
存储 分布式计算 Apache
万字长文:基于Apache Hudi + Flink多流拼接(大宽表)最佳实践
万字长文:基于Apache Hudi + Flink多流拼接(大宽表)最佳实践
122 3
|
5月前
|
数据采集 机器学习/深度学习 数据可视化
使用Python进行数据分析的最佳实践
数据分析已经成为了现代生活和商业决策中的不可或缺的一部分。Python是数据分析的首选编程语言之一,因为它具有丰富的库和工具,可以轻松处理、可视化和分析数据。本文将探讨使用Python进行数据分析的最佳实践,帮助你提高工作效率和数据分析的质量。
|
9月前
|
SQL 存储 消息中间件
Flink+StarRocks 实时数据分析新范式
StarRocks 社区技术布道师谢寅,在 Flink Forward Asia 2022 实时湖仓的分享。
1367 2
Flink+StarRocks 实时数据分析新范式
|
10月前
|
消息中间件 SQL Cloud Native
[实战系列]SelectDB Cloud Flink Connector 最佳实践
随着云基础设施的不断完善,云原生已经成为各行业数字化转型的必选项,越来越多的应用开始进行云原生化架构升级和应用迁移。 而云原生实时数仓的出现,让传统的数据仓库无论是成本、灵活性还是开放性等方面都显露出不足。拥有高性能、高可用性、可伸缩性、高安全性等特征的云原生数据库,正在成为企业的首选。 SelectDB Cloud作为一款运行于多云之上的云原生实时数据仓库,可以为客户提供极简运维和极致性价比的数仓服务,为用户提供开箱即用的能力。 同时,SelectDB Cloud 结合 Flink 流式计算,可以让用户将 Kafka 中的非结构化数据以及 MySQL 等上游业务库中的变更数据,实时同步到 S
120 0
|
11月前
|
分布式计算 监控 前端开发
《Apache Flink 案例集(2022版)》——2.数据分析——网易互娱-基于Flink 的支付环境全关联分析实践(上)
《Apache Flink 案例集(2022版)》——2.数据分析——网易互娱-基于Flink 的支付环境全关联分析实践(上)
135 0
|
关系型数据库 MySQL Java
为什么 Flink 无法实时写入 MySQL?
Flink 1.10 使用 flink-jdbc 连接器的方式与 MySQL 交互,读数据和写数据都能完成,但是在写数据时,发现 Flink 程序执行完毕之后,才能在 MySQL 中查询到插入的数据。即,虽然是流计算,但却不能实时的输出计算结果?
为什么 Flink 无法实时写入 MySQL?
|
2月前
|
消息中间件 Kafka Apache
Apache Flink 是一个开源的分布式流处理框架
Apache Flink 是一个开源的分布式流处理框架
482 5
|
1月前
|
SQL Java API
官宣|Apache Flink 1.19 发布公告
Apache Flink PMC(项目管理委员)很高兴地宣布发布 Apache Flink 1.19.0。
1341 1
官宣|Apache Flink 1.19 发布公告

相关产品

  • 实时计算 Flink版