阿里云实时计算Flink版的优势极大:

  • 性能优越:作业可达百万级吞吐,计算可达秒级延迟,TPC-H性能测试可达开源引擎3~5倍。
  • 功能强大:数十种作业指标监控,一站式开发界面,提供智能诊断系统,具有作业智能调优功能。
  • 价格低廉:极致弹性体验,可按量付费,总资源费用低于自建。
  • 稳定安全:服务SLA可达99.9%,集群计算无单点,故障可自动恢复,资源租户隔离,杜绝相互干扰。
  • 品牌认证:Flink官方创始团队出品,中国信通院认证,进入Forrester象限的实时流计算产品。
  • 兼容开源:提供最新Flink版本,与开源Flink接口100%兼容,实现业务平滑迁移上云。

这边广告费该找谁结一下,,,

正文开始之前先说一下主要的属性

数据来源:

  1. 实时表:datahub
  2. 代码表:rds(阿里云的云数据库,这里以MySQL为证)
  3. 结果表: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;

 

Logo

魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。

更多推荐