Flink测试利器之DataGen初探 | 京东云技术团队

什么是 Flinksql

Flink SQL 是基于 Apache Calcite 的 SQL 解析器和优化器构建的,支持ANSI SQL 标准,允许使用标准的 SQL 语句来处理流式和批处理数据。通过 Flink SQL,可以以声明式的方式描述数据处理逻辑,而无需编写显式的代码。使用 Flink SQL,可以执行各种数据操作,如过滤、聚合、连接和转换等。它还提供了窗口操作、时间处理和复杂事件处理等功能,以满足流式数据处理的需求。

Flink SQL 提供了许多扩展功能和语法,以适应 Flink 的流式和批处理引擎的特性。他是Flink最高级别的抽象,可以与 DataStream API 和 DataSet API 无缝集成,利用 Flink 的分布式计算能力和容错机制。

使用 Flink SQL处理数据的基本步骤:

  1. 定义输入表:使用 CREATE TABLE 语句定义输入表,指定表的模式(字段和类型)和数据源(如 Kafka、文件等)。

  2. 执行 SQL 查询:使用 SELECT、INSERT INTO 等 SQL 语句来执行数据查询和操作。您可以在 SQL 查询中使用各种内置函数、聚合操作、窗口操作和时间属性等。

  3. 定义输出表:使用 CREATE TABLE 语句定义输出表,指定表的模式和目标数据存储(如 Kafka、文件等)。

  4. 提交作业:将 Flink SQL 查询作为 Flink 作业提交到 Flink 集群中执行。Flink会根据查询的逻辑和配置自动构建执行计划,并将数据处理任务分发到集群中的任务管理器进行执行。

总而言之,我们可以通过Flink SQL 查询和操作来处理流式和批处理数据。它提供了一种简化和加速数据处理开发的方式,尤其适用于熟悉 SQL 的开发人员和数据工程师。

什么是 connector

Flink Connector 是指用于连接外部系统和数据源的组件。它允许 Flink 通过特定的连接器与不同的数据源进行交互,例如数据库、消息队列、文件系统等。它负责处理与外部系统的通信、数据格式转换、数据读取和写入等任务。无论是作为输入数据表还是输出数据表,通过使用适当的连接器,可以在 Flink SQL 中访问和操作外部系统中的数据。目前实时平台提供了很多常用的连接器:

例如:

  1. JDBC :用于与关系型数据库(如 MySQL、PostgreSQL)建立连接,并支持在 Flink SQL 中读取和写入数据库表的数据。

  2. JDQ :用于与 JDQ 集成,可以读取和写入 JDQ 主题中的数据。

  3. Elasticsearch :用于与 Elasticsearch 集成,可以将数据写入 Elasticsearch 索引或从索引中读取数据。

  4. File Connector:用于读取和写入各种文件格式(如 CSV、JSON、Parquet)的数据。

  5. ......

还有如HBase、JMQ4、Doris、Clickhouse,Jimdb,Hive等,用于与不同的数据源进行集成。通过使用 Flink SQL Connector,我们可以轻松地与外部系统进行数据交互,将数据导入到 Flink 进行处理,或将处理结果导出到外部系统。

DataGen Connector

DataGen 是 Flink SQL 提供的一个内置连接器,用于生成模拟的测试数据,以便在开发和测试过程中使用。

使用 DataGen,可以生成具有不同数据类型和分布的数据,例如整数、字符串、日期等。这样可以模拟真实的数据场景,并帮助验证和调试 Flink SQL 查询和操作。

demo

以下是一个使用 DataGen 函数的简单示例:

1-- 创建输入表 2CREATE TABLE input_table ( 3 order_number BIGINT, 4 price DECIMAL(32,2), 5 buyer ROW<first_name STRING, last_name STRING>, 6 order_time TIMESTAMP(3) 7) WITH ( 8 'connector' = 'datagen', 9); 10

在上面的示例中,我们使用 DataGen 连接器创建了一个名为 `input_table` 的输入表。该表包含了 `order_number`、`price` 和 `buyer` ,`order_time`四个字段。默认是random随机生成对应类型的数据,生产速率是10000条/秒,只要任务不停,就会源源不断的生产数据。当然也可以指定一些参数来定义生成数据的规则,例如每秒生成的行数、字段的数据类型和分布。

生成的数据样例:

1{"order_number":-6353089831284155505,"price":253422671148527900374700392448,"buyer":{"first_name":"6e4df4455bed12c8ad74f03471e5d8e3141d7977bcc5bef88a57102dac71ac9a9dbef00f406ce9bddaf3741f37330e5fb9d2","last_name":"d7d8a39e063fbd2beac91c791dc1024e2b1f0857b85990fbb5c4eac32445951aad0a2bcffd3a29b2a08b057a0b31aa689ed7"},"order_time":"2023-09-21 06:22:29.618"} 2{"order_number":1102733628546646982,"price":628524591222898424803263250432,"buyer":{"first_name":"4738f237436b70c80e504b95f0d9ec3d7c01c8745edf21495f17bb4d7044b4950943014f26b5d7fdaed10db37a632849b96c","last_name":"7f9dbdbed581b687989665b97c09dec1a617c830c048446bf31c746898e1bccfe21a5969ee174a1d69845be7163b5e375a09"},"order_time":"2023-09-21 06:23:01.69"} 3

支持的类型

字段类型数据生成方式
BOOLEANrandom
CHARrandom / sequence
VARCHARrandom / sequence
STRINGrandom / sequence
DECIMALrandom / sequence
TINYINTrandom / sequence
SMALLINTrandom / sequence
INTrandom / sequence
BIGINTrandom / sequence
FLOATrandom / sequence
DOUBLErandom / sequence
DATErandom
TIMErandom
TIMESTAMPrandom
TIMESTAMP_LTZrandom
INTERVAL YEAR TO MONTHrandom
INTERVAL DAY TO MONTHrandom
ROWrandom
ARRAYrandom
MAPrandom
MULTISETrandom

连接器属性

属性是否必填默认值类型描述
connectorrequired(none)String'datagen'.
rows-per-secondoptional10000Long数据生产速率
number-of-rowsoptional(none)Long指定生产的数据条数,默认是不限制。
fields.#.kindoptionalrandomString指定字段的生产数据的方式 random还是sequence
fields.#.minoptional(Minimum value of type)(Type of field)random生成器 指定字段 # 最小值, 支持数字类型
fields.#.maxoptional(Maximum value of type)(Type of field)random生成器的指定字段 # 最大值, 支持数字类型
fields.#.lengthoptional100Integerchar/varchar/string/array/map/multiset 类型的长度.
fields.#.startoptional(none)(Type of field)sequence生成器的开始值
fields.#.endoptional(none)(Type of field)sequence生成器的结束值

DataGen使用

了解了dategen的基本使用方法,那么下面来结合其他类型的连接器实践一下吧。

场景1 生成一亿条数据到hive表

1CREATE TABLE dataGenSourceTable 2 ( 3 order_number BIGINT, 4 price DECIMAL(10, 2), 5 buyer STRING, 6 order_time TIMESTAMP(3) 7 ) 8WITH 9 ( 'connector'='datagen', 10 'number-of-rows'='100000000', 11 'rows-per-second' = '100000' 12 ) ; 13 14 15CREATECATALOG myhive 16WITH ( 17 'type'='hive', 18 'default-database'='default' 19); 20USECATALOG myhive; 21USE dev; 22SETtable.sql-dialect=hive; 23CREATETABLEifnotexists shipu3_test_0932 ( 24 order_number BIGINT, 25 price DECIMAL(10, 2), 26 buyer STRING, 27 order_time TIMESTAMP(3) 28) PARTITIONED BY (dt STRING) STORED AS parquet TBLPROPERTIES ( 29 'partition.time-extractor.timestamp-pattern'='$dt', 30 'sink.partition-commit.trigger'='partition-time', 31 'sink.partition-commit.delay'='1 h', 32 'sink.partition-commit.policy.kind'='metastore,success-file' 33); 34SETtable.sql-dialect=default; 35insert into myhive.dev.shipu3_test_0932 36select order_number,price,buyer,order_time, cast( CURRENT_DATE as varchar) 37from default_catalog.default_database.dataGenSourceTable; 38

当每秒生产10万条数据的时候,17分钟左右就可以完成,当然我们可以通过增加Flink任务的计算节点、并行度、提高生产速率'rows-per-second'的值等来更快速的完成大数据量的生产。

场景2 持续每秒生产10万条数到消息队列

1CREATE TABLE dataGenSourceTable ( 2 order_number BIGINT, 3 price INT, 4 buyer ROW< first_name STRING, last_name STRING >, 5 order_time TIMESTAMP(3), 6 col_array ARRAY < STRING >, 7 col_map map < STRING, STRING > 8 ) 9WITH 10 ( 'connector'='datagen', --连接器类型 11 'rows-per-second'='100000', --生产速率 12 'fields.order_number.kind'='random', --字段order_number的生产方式 13 'fields.order_number.min'='1', --字段order_number最小值 14 'fields.order_number.max'='1000', --字段order_number最大值 15 'fields.price.kind'='sequence', --字段price的生产方式 16 'fields.price.start'='1', --字段price开始值 17 'fields.price.end'='1000', --字段price最大值 18 'fields.col_array.element.length'='5', --每个元素的长度 19 'fields.col_map.key.length'='5', --map key的长度 20 'fields.col_map.value.length'='5' --map value的长度 21 ) ; 22CREATE TABLE jdqsink1 23 ( 24 order_number BIGINT, 25 price DECIMAL(32, 2), 26 buyer ROW< first_name STRING, last_name STRING >, 27 order_time TIMESTAMP(3), 28 col_ARRAY ARRAY < STRING >, 29 col_map map < STRING, STRING > 30 ) 31WITH 32 ( 33 'connector'='jdq', 34 'topic'='jrdw-fk-area_info__1', 35 'jdq.client.id'='xxxxx', 36 'jdq.password'='xxxxxxx', 37 'jdq.domain'='db.test.group.com', 38 'format'='json' 39 ) ; 40INSERTINTO jdqsink1 41SELECT*FROM dataGenSourceTable; 42

思考

通过以上案例可以看到,通过Datagen结合其他连接器可以模拟各种场景的数据

  • 性能测试:我们可以利用Flink的高处理性能,来调试任务的外部依赖的阈值(超时,限流等)到一个合适的水位,避免自己的任务有过多的外部依赖出现木桶效应;
  • 边界条件测试:我们通过使用 Flink DataGen 生成特殊的测试数据,如最小值、最大值、空值、重复值等来验证 Flink 任务在边界条件下的正确性和鲁棒性;
  • 数据完整性测试:我们通过Flink DataGen 可以生成包含错误或异常数据的数据集,如无效的数据格式、缺失的字段、重复的数据等。从而可以测试 Flink 任务对异常情况的处理能力,验证 Flink任务在处理数据时是否能够正确地保持数据的完整性。

总之,Flink DataGen 是一个强大的工具,可以帮助测试人员构造各种类型的测试数据。通过合理的使用 ,测试人员可以更有效地进行测试,并发现潜在的问题和缺陷。

作者:京东零售 石朴

来源:京东云开发者社区 转载请注明来源

点赞
收藏

评论区

加载中...

相关推荐

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

mysql中like用法

like的通配符有两种%(百分号):代表零个、一个或者多个字符。\(下划线):代表一个数字或者字符。1\.name以"李"开头wherenamelike'李%'2\.name中包含"云",“云”可以在任何位置wherenamelike'%云%'3\.第二个和第三个字符是0的值wheresalarylike'\00%'4\

FLV文件格式

1.        FLV文件对齐方式FLV文件以大端对齐方式存放多字节整型。如存放数字无符号16位的数字300(0x012C),那么在FLV文件中存放的顺序是:|0x01|0x2C|。如果是无符号32位数字300(0x0000012C),那么在FLV文件中的存放顺序是:|0x00|0x00|0x00|0x01|0x2C。2.  

SpringBoot整合Redis乱码原因及解决方案

问题描述:springboot使用springdataredis存储数据时乱码rediskey/value出现\\xAC\\xED\\x00\\x05t\\x00\\x05问题分析:查看RedisTemplate类!(https://oscimg.oschina.net/oscnet/0a85565fa

mysql设置时区

mysql设置时区mysql\_query("SETtime\_zone'8:00'")ordie('时区设置失败,请联系管理员!');中国在东8区所以加8方法二:selectcount(user\_id)asdevice,CONVERT\_TZ(FROM\_UNIXTIME(reg\_time),'08:00','0