阿里云实时计算平台Flink的作业开发流程详解
·
阿里云实时计算Flink版的优势极大:
- 性能优越:作业可达百万级吞吐,计算可达秒级延迟,TPC-H性能测试可达开源引擎3~5倍。
- 功能强大:数十种作业指标监控,一站式开发界面,提供智能诊断系统,具有作业智能调优功能。
- 价格低廉:极致弹性体验,可按量付费,总资源费用低于自建。
- 稳定安全:服务SLA可达99.9%,集群计算无单点,故障可自动恢复,资源租户隔离,杜绝相互干扰。
- 品牌认证:Flink官方创始团队出品,中国信通院认证,进入Forrester象限的实时流计算产品。
- 兼容开源:提供最新Flink版本,与开源Flink接口100%兼容,实现业务平滑迁移上云。
(这边广告费该找谁结一下,,,)
正文开始之前先说一下主要的属性
数据来源:
- 实时表:datahub
- 代码表:rds(阿里云的云数据库,这里以MySQL为证)
- 结果表:rds
开发思路:把datahub里实时采集的数据进行实时计算并把结果实时写入到表
开发流程:建立与datahub相关联对应的结构一致的数据表,建立与rds代码表结构一致的代码表,建立与rds结果表结构一致的结果表。这里为什么建立三张不同属性的表呢?因为不论是数据来源还是数据写入的目的地,所有的数据都不是实时计算平台的主要存储地,平台只是提供一个接收与计算的临时站,所以需要在平台建立对应的结构一致的表来临时接收和储存数据,需要明确的在同一个开发作业里,业务表,代码表,结果表都可以是多个的,但这样会影响作业的并行度,建议一个作业一个结果表,并数据开发写入结果表。
废话不多说直接上作业开发:
- 首先新建job作业,然后在开发框里写上注释,以下所有的代码均可写在同一个作业里
--Blink SQL
--********************************************************************--
--Author : Mochou
--CreateTime: 2020-10-19 16:50:17
--Comment : 实时展示动物微笑的累计次数
--********************************************************************--
- 创建topic的同名业务表,以把topic里的数据实时拉到此表
CREATE TABLE t_animal (
ID VARCHAR,
NUM1 INTEGER,
SMILESJ DATE
) WITH (
type = 'datahub',
endPoint = 'http://******************',
roleArn = 'acs:ram::*******************',
project = '${projectName}',
topic = '${topicName}'
);
- 创建rds代码表
CREATE TABLE t_dim (
ID VARCHAR,
NAME VARCHAR
PRIMARY KEY (ID),
-- 下面代码的意思是表数据依据系统时间
PERIOD FOR SYSTEM_TIME
) WITH (
type = 'rds',
url = 'jdbc:mysql://*******************',
userName = '${userName}',
password = '${password}',
tableName = 't_dim'
);
- 创建rds汇总结果表
CREATE TABLE t_animal_res (
ID VARCHAR,
NAME VARCHAR,
NUM1 INTEGER,
RQ VARCHAR
PRIMARY KEY (ID)
) WITH (
type = 'rds',
url = 'jdbc:mysql://*******************',
userName = '${userName}',
password = '${password}',
tableName = 't_animal_res'
);
- 数据开发,把两张表的数据实时写入到上述结果表
INSERT INTO
t_animal_res
SELECT
B.ID,
B.NAME,
SUM(A.NUM1) AS NUM1,
SUBSTR (CAST (CURRENT_TIMESTAMP AS VARCHAR), 1, 10) AS RQ
FROM
t_animal A
JOIN t_dim FOR SYSTEM_TIME AS OF PROCTIME () AS B ON A.ID = B.ID
WHERE
--按照数据同步时间进行where关联当天
SUBSTR (CAST (A.SMILESJ AS VARCHAR), 1, 10) = SUBSTR (CAST (CURRENT_TIMESTAMP AS VARCHAR), 1, 10)
GROUP BY
B.ID,
B.NAME;
全代码如下:
--Blink SQL
--********************************************************************--
--Author : Mochou
--CreateTime: 2020-10-19 16:50:17
--Comment : 实时展示动物微笑的累计次数
--********************************************************************--
CREATE TABLE t_animal (
ID VARCHAR,
NUM1 INTEGER,
SMILESJ DATE
) WITH (
type = 'datahub',
endPoint = 'http://******************',
roleArn = 'acs:ram::*******************',
project = '${projectName}',
topic = '${topicName}'
);
CREATE TABLE t_dim (
ID VARCHAR,
NAME VARCHAR
PRIMARY KEY (ID),
PERIOD FOR SYSTEM_TIME
) WITH (
type = 'rds',
url = 'jdbc:mysql://*******************',
userName = '${userName}',
password = '${password}',
tableName = 't_dim'
);
CREATE TABLE t_animal_res (
ID VARCHAR,
NAME VARCHAR,
NUM1 INTEGER,
RQ VARCHAR
PRIMARY KEY (ID)
) WITH (
type = 'rds',
url = 'jdbc:mysql://*******************',
userName = '${userName}',
password = '${password}',
tableName = 't_animal_res'
);
INSERT INTO
t_animal_res
SELECT
B.ID,
B.NAME,
SUM(A.NUM1) AS NUM1,
SUBSTR (CAST (CURRENT_TIMESTAMP AS VARCHAR), 1, 10) AS RQ
FROM
t_animal A
JOIN t_dim FOR SYSTEM_TIME AS OF PROCTIME () AS B ON A.ID = B.ID
WHERE
--按照数据同步时间进行where关联当天
SUBSTR (CAST (A.SMILESJ AS VARCHAR), 1, 10) = SUBSTR (CAST (CURRENT_TIMESTAMP AS VARCHAR), 1, 10)
GROUP BY
B.ID,
B.NAME;
魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。
更多推荐


所有评论(0)