业务流程组件配置参数介绍

配置各类组件参数的步骤如下:

  1. 将组件拖拽到画布中后,双击该组件,右侧边栏会弹出编辑窗口。

  1. 配置组件的相关参数,参数说明请参见后续各组件介绍中的表格。

  1. 单击<确定>按钮,组件配置完成。

数据集成组件

数据集成组件是业务流程画布中用来关联数据集成作业的组件,包含同步DI、异步DI和数据同步三种组件。

数据集成的作用是通过集成服务,将本系统外的数据经过一定的处理后存入本系统相关的数据存储系统中。

表-1 数据集成组件配置参数介绍

组件

说明

同步DI、异步DI

DI组件用来关联[融合集成平台/数据集成/作业管理]模块下的DI作业。关于DI作业和新增方法的介绍可参见[融合集成平台/数据集成/作业管理]的联机帮助。该组件配置参数说明如下:

  • 节点名称:配置数据集成节点的名称

  • DI作业:选择DI作业

节点类型不可编辑;其他参数会根据DI作业自动加载,不可编辑

数据同步

数据同步组件用来关联数据同步任务。关于数据同步任务的介绍和新建,请参见数据同步管理。该组件配置参数说明如下:

  • 节点名称:配置数据同步节点的名称

  • 任务名称:选择当前组织下未被其他业务流程使用的数据同步任务

  • KAFKA数据源:Kafka数据源,选择任务名称后会自动填充

  • 管道名称:管道,选择任务名称后会自动填充

  • 任务类型:任务的类型,选择任务名称后会自动填充

  • 创建时间:任务的创建时间,选择任务名称后会自动填充

  • 执行队列:选择本组织可以使用的队列

  • 并行度:根据实际需求自定义设置并行度

节点类型不可编辑;KAFKA数据源、管道名称、任务类型和创建时间参数会根据选择的数据同步任务自动加载,不可编辑

 

离线分析组件

如离线分析相关的任务涉及对STDB进行查询等操作,需在[配置管理/参数管理]中,将“离线分析是否加载时空函数”配置为true

 

离线分析组件是业务流程画布中用来关联离线作业的组件,包含MRSparkJarHiveSQLJavaSqoopShellSparkSQLPySparkRDSSQL类型。

离线分析

说明

MRSparkJarJavaShellPySpark

可以选择任务管理下对应类型的类型的任务作为模板,加载后适当修改参数,并配置执行队列,即可为完成配置(除执行队列外的其他参数是非必要的)

执行队列:选择本组织可以使用的队列

HiveSQL

可以配置作业名称、执行队列、SQLJDBC URL,即可完成(其他参数是非必要的)。该组件配置参数说明如下:

  • 节点名称:配置作业节点的名称

  • 作业名称:必填,配置作业的名称

  • 描述:配置作业的描述信息

  • 执行队列:必选,选择本组织可以使用的队列

  • prepare:对应Oozie中的prepare标签元素,用于创建或者删除指定的HDFS目录

  • configuration:指定Mapper类和Reducer类等MapReduce客户端配置参数,鼠标悬于后方的图标上,可以查看该参数的说明。支持配置多项,点击后方的图标可以新增一项,点击图标可以删除一项

  • sql:必填,单击<编辑SQL>按钮,在弹出窗口中输入SQL语句,输入完成后单击<确定>按钮

  • jdbcUrl:必填,配置连接HiveJDBC URL

  • passWord:配置连接Hive所需的密码

  • param:指定执行所需的环境变量。支持配置多条,点击后方的图标可以新增一条,点击图标可以删除一条

  • 出错重跑:配置是否在运行失败的情况下重新运行

  • 出错重跑次数:配置运行失败时的重新运行次数,开启出错重跑后需配置该参数

  • 出错重跑间隔:配置失败后重新运行的时间间隔,开启出错重跑后需配置该参数

Sqoop

  • 节点名称:配置作业节点的名称

  • 作业名称:必填,配置作业的名称

  • 描述:配置作业的描述信息

  • 执行队列:必选,选择本组织可以使用的队列

  • prepare:对应Oozie中的prepare标签元素,用于创建或者删除指定的HDFS目录

  • configuration:指定Mapper类和Reducer类等MapReduce客户端配置参数。支持配置多项,点击后方的图标可以新增一项,点击图标可以删除一项

  • 功能:选择Sqoop组件的功能,即从HDFSRDB转换,从RDBHDFS转换,以及手动配置命令等

  • RDB:在弹出窗口中配置RDB类型数据源中的表名

  • HDFS:在弹出窗口中选择HDFS目录

  • COMMAND:手动配置命令

  • 出错重跑:配置是否在运行失败的情况下重新运行

  • 出错重跑次数:配置运行失败时的重新运行次数,开启出错重跑后需配置该参数

  • 出错重跑间隔:配置失败后重新运行的时间间隔,开启出错重跑后需配置该参数

SparkSQL

是用来关联需要到Hadoop集群上以SparkSQL方式运行的Spark任务。其参数说明如下:

  • 节点名称:配置作业节点的名称

  • 描述:配置作业的描述信息

  • 执行队列:必选,指定SparkSQL作业在集群中运行使用的队列

  • sql:单击<编辑SQL>按钮,可以在弹出窗口中编辑SQL语句。编辑SQL语句的操作介绍请参见编辑SQL

  • executor核数:常用的SparkSQL作业节点的规格参数。配置Spark 单个executor执行需要的cores个数

  • 内存(MB):常用的SparkSQL作业节点的规格参数。配置Spark 单个executor执行需要的内存大小

  • driver核数:常用的SparkSQL作业节点的规格参数。配置Spark driver程序执行需要的cores个数

  • 内存(MB):常用的SparkSQL作业节点的规格参数。配置Spark driver程序执行需要的内存大小

  • 实例数:常用的SparkSQL作业节点的规格参数。配置Spark任务执行需要的 executor的实例个数

  • 堆外内存:常用的SparkSQL作业节点的规格参数。配置Spark任务的executor程序堆外内存大小

  • driver堆外内存:常用的SparkSQL作业节点的规格参数。配置Spark dirver程序执行需要的堆外内存大小

  • sparkOpts:配置SparkSQL作业节点的规格参数。输入框中存在缺省的配置值,可以根据实际需求调整。如运行SparkSQL还需要其他Spark相关的配置项均可在此处点击添加。需要删除则点击。缺省值的介绍如下:

  • spark.sql.autoBroadcastJoinThreshold:当执行join时,被广播到worker节点上表最大字节。默认10MB,配置自动广播的数据量大小。如果设置成-1,则禁用该特性

  • spark.dynamicAllocation.maxExecutors:动态分配最大executor个数,默认设置30个。可根据实际需求调整

  • 参数配置:配置作业的时间范围参数

  • configuration:指定Mapper类和Reducer类等MapReduce客户端配置参数。支持配置多项,点击后方的图标可以新增一项,点击图标可以删除一项

  • 出错重跑:配置是否在运行失败的情况下重新运行

  • 出错重跑次数:配置运行失败时的重新运行次数,开启出错重跑后需配置该参数

  • 出错重跑间隔:配置失败后重新运行的时间间隔,开启出错重跑后需配置该参数

RDSSQL

下发分析型SQL语句到MPP(GP,VERTICA)数据库,其参数说明如下:

  • 节点名称:配置作业节点的名称

  • 描述:配置作业的描述信息

  • 执行队列:必选,选择本组织可以使用的YARN队列

  • sql:单击<编辑SQL>按钮,可以在弹出窗口中编辑SQL语句。编辑SQL语句的操作介绍请参见RDSSQL编辑SQL

  • 出错重跑:配置是否在运行失败的情况下重新运行

  • 出错重跑次数:配置运行失败时的重新运行次数,开启出错重跑后需配置该参数

  • 出错重跑间隔:配置失败后重新运行的时间间隔,开启出错重跑后需配置该参数

 

对于SparkSQL组件,其计算结果可能会出现精度有误情况,如-3所示。

表-2 运行结果中精度有误类型列表

数据源表类型

精度有误类型

PostgreSQL

timebitreal

Greenplum

decimalnumericrealtimebit

达梦

timetimestampbitbinary

MySQL

timeyearbitbooleantinyint

Oracle

numberfloat

Vertica

timetimetzgeographygeometry

 

MySQLtinyint(1)类型的字段在注册离线表后,非01的值将被转换为true,造成数据失真。

 

 

 

编辑SQL

为优化集群资源参数,在集群开启Kerberos认证的情况下,请至[系统]模块下的[集群管理/集群资源]配置页面中,为所在组织配置Kerberos用户及队列信息。

 

SparkSQL组件、HiveSQL组件或RDSSQL组件的编辑SQL功能提供了SQL语句编辑、执行和校验等功能。在弹出的SQL窗口中:

当前SQL编辑器中对查询返回的数据条数进行了限制,查询结果最大返回1000条数据。

如果对数据量过大、逻辑复杂的SQL进行了语法校验,可能会导致超时或者执行失败,并影响Spark Thrift Server的稳定性。

 

1.      执行SQL

SQL语句编辑完成后,可以执行SQL并查看执行结果,以确认SQL语句是否符合作业要求。

SQL编辑区域中提供了对选中的SQL语句和完整SQL语句的执行的功能。执行选中SQL时,仅支持管理员用户进行删除操作。查询结果默认返回最多1000条结果。

2.      SQL上传

单击编辑区域上方的<SQL上传>按钮,选择已经编辑好的SQL语句文件,可以快速生成SQL语句。

3.      格式化

单击编辑区域上方的<格式化>按钮,可以将编辑区域的SQL语句格式进行标准化,以方便查看。

4.      语法校验

SQL语句编辑完成后,可以校验SQL语句语法,以检查SQL语句是否符合语法要求。

SQL编辑区域中提供了对选中的SQL语句和完整SQL语句的语法校验功能。

5.      结果导出

对于SparkSQL,还支持配置是否导出到其他数据源中存储。如果选择导出,则需要配置导出参数。

表-3 结果导出参数说明

参数

说明

导出参数

  • 导出类型:选择导出的类型,包括HDFSTable。选择不同的导出类型,需要配置的参数不同。

  • HDFS目录:选择HDFS的目录(如果选择HDFS,才需配置该参数)。点击输入框,弹出HDFS目录窗口,在窗口的列表中选择目标文件夹,并单击<确定>按钮即可。窗口中提供了图标,可在当前目录下新增子目录;提供了图标,可以返回上层目录;提供了图标,可以刷新窗口;列表中还提供了文件夹和文件删除操作链接,以及文件的下载链接。

  • 数据分隔符:选择分隔数据的标识符(如果选择HDFS,才需配置该参数)

  • 表类型:选择表所属数据源的类型(如果选择Table,才需配置该参数)

  • 数据源:选择数据源(如果选择Table,才需配置该参数)

  • schema:选择schema(如果选择Table,且“表类型”参数选择PostgreSQLGreenplum、达梦、Vertica时,才需配置该参数)

  • 表名:选择或者输入表名称(如果选择Table,才需配置该参数)

  • 写入模式:选择写入目录的模式。覆盖模式会覆盖掉原有的内容,请谨慎选择。为保证作业可以直接多次运行,新建模式的作业运行时会被适配为追加模式

  • SparkSQL组件导出到HDFS的写入方式有追加、覆盖。追加指追加查询数据文件到所选的HDFS目录下;覆盖指清空选择的HDFS目录并将查询结果导出到所选的HDFS目录下

  • SparkSQL组件导出到Table的写入方式有新建、追加、覆盖、更新(仅Vertica表类型支持)。导出的映射关系要求为:SQL计算结果表的列名须与导出表列名一致,数据类型需与导出表数据类型匹配。新建指新建表并将查询数据导出到该表中;追加指追加查询数据文件到所选的表中,如该表不存在,则会首先新建表;覆盖指清空所选的表并将查询结果导出到该表中;更新指根据用户所选更新条件查询导出表中数据,有则更新,无则插入

 

 

 

实时计算组件

实时计算组件是业务流程画布中用来关联实时作业的组件,包含StreamingJob组件。实时计算组件可以选取实时作业管理下的实时作业。

StreamingJob组件是用来关联需要到Hadoop集群上运行的Flink任务。

该组件配置参数说明如下:

节点类型不可编辑;其他参数会根据实时作业任务自动加载,不可编辑。

可根据需要配置作业的运行资源,包括执行队列、管理节点内存、工作节点内存、作业的并行度(作业并行度最小值为1,最大值分两种情况:当数据运营平台集群未配置时,并行度最大值为30;当数据运营平台集群配置后,并行度最大值为所使用的DataEngine大数据集群中Kafka组件节点数量的10倍)以及检查点的间隔时间和模式。

对于管理节点内存和工作节点内存配置参数,请根据作业所要处理的数据量以及机器资源情况,尽可能调大两个参数的值,避免作业因为所要处理的数据量过大导致作业内存溢出等情况的发生。

对于作业并行度参数,建议用户根据作业所要使用的数据源表情况(Kafka表对应的Topic分区数量)来设置该参数。

高级参数对于FLINK_JAR类型的实时作业为默认参数,若用户在作业内部(Jar文件中的代码内)设置了以上参数,则以作业内部设置的参数为准。

 

控制节点组件

控制节点组件是用于对业务中作业的运行进行控制的组件,包括Fork组件、Join组件、跨组织节点组件和分支节点组件。

表-4 控制节点组件介绍

组件

说明

ForkJoin

Fork组件和Join组件自身不会参与任何作业的数据处理,且Fork组件与Join组件在业务流程中需成对出现。

  • Fork组件:业务流程中某作业的后续作业需要并行执行时,需添加该组件以连接前序作业和后续的并行作业。

  • Join组件:业务流程中某几个并行作业需全部完成后才执行后续的作业,需添加该组件以连接前序的并行作业和后续作业。

例如,在业务流程画布中,当某个作业A的后续为两个需要并行执行的作业B和作业C,且两个作业均执行完成后才可继续执行后续作业D时,需要在作业A后添加Fork组件节点,然后从Fork组件节点连接至作业B和作业C,之后将作业B与作业C连接至Join组件节点,再将Join组件连接至作业D

跨组织节点

跨组织节点提供了事件触发调度的能力,包含发布、接收两种子类型,分别承担消息发布和接收处理的职责。通过引入跨组织节点,可以解决存在依赖关系的业务流程的启动顺序问题。

  • 跨组织(发布)节点:运行离线作业,发布消息,标识该节点运行结束。

  • 跨组织(接收)节点:配置监听任务,当接收到所有监听目标发布的消息后,启动本作业。

例如,当前有两个业务流程AB,流程B依赖流程A的分析结果,因此必须在流程A运行结束后,再启动流程B,若仅依赖手动启动,会存在启动不及时且使用不便的问题。此时可以通过跨组织节点来解决,在流程A的末尾连接跨组织(发布)节点,在流程B的起始位置连接跨组织(接收)节点,并配置流程B的接收节点监听流程A的发布节点,流程A运行结束后会发布消息,流程B的接收节点会接收并处理消息,最终启动流程B

分支节点

Fork-Join一样,分支(Decision)节点也是数据开发提供的逻辑控制类节点中的一种。通过定义分支条件,在实际运行时可以根据分支条件的内容是否为true,决定分支下游的走向。

例如,在业务流程画布中,当某个作业A的后续为两个需要二选一执行的作业B和作业C。此时需要在A作业之后添加分支节点,将分支节点与下游BC节点相连接,并配置分支节点的属性,即分别对B节点和C节点配置执行的分支条件。当两个条件中任何一个结果为true,则执行对应的下游分支节点;当两个条件均满足条件时,选择第一个节点执行;当两个节点均不满足条件时,整个业务流程的执行会跳到最后直接结束,后续节点均不再执行。

 

跨组织节点说明:

·          DI作业内置了跨组织(发布)节点的能力,任务结束时,发布消息。

·          跨组织(发布)节点可以不被任何节点监听,独立存在。

·          跨组织(接收)节点支持两种监听类型,分别是跨组织(发布)节点和DI作业。

·          包含跨组织(接收)节点的业务流程,仅响应流程启动后监听目标发布的消息,忽略历史消息。

·          跨组织(接收)节点启动业务流程时,会放弃提交处于运行状态的作业。此时该条消息已经被消费,任务结束后等待后续消息进行调度。

跨组织(发布)节点存在如下使用限制:

·          必须与其他离线作业成组使用

·          同组内有且仅有一个父节点

·          父节点不能为跨组织(发布)节点

跨组织(接收)节点存在如下使用限制:

·          必须与其他离线作业成组使用

·          必须是组内的起始节点

·          同组内有且仅有一个子节点