通过消息投递功能,您可以将目标项目空间下Flink作业的启动日志、资源使用情况、运行事件及作业状态变更消息实时传输至日志服务SLS或云消息队列Kafka版。该功能可以满足数据持久化、历史记录保存、灵活集成和实时分析的需求,同时便于随时查询历史数据,为问题排查、性能优化及审计分析提供支持。本文为您介绍消息投递的操作与查看方法。
背景信息
消息投递支持投递至日志服务SLS和云消息队列Kafka版,两个渠道的消息内容一致。SLS支持跨地域投递,Kafka仅支持投递至与Flink工作空间同地域的实例。目前支持投递信息的类型及投递时机详情如下。
|
类别 |
投递说明 |
投递时机 |
|
作业启动日志 |
投递作业启动后的日志,涵盖从Flink环境被初始化到Job Manager启动,并生成Flink运行图的整个启动过程所产生的日志。 |
作业启动成功或者作业达到终态(失败或结束), 投递一次启动日志。 |
|
资源情况 |
重要
仅用于资源容量管理的追踪,不具备告警功能。 |
项目空间正常运行情况下,每30秒发送一次资源情况。 |
|
运行事件 |
投递作业启动的过程中每个时间点的启动状态。 |
作业运行事件生成, 触发投递。 |
|
作业消耗资源 |
获取运行中流作业的资源使用情况。其中不包括批作业和Session集群运行的作业。 |
作业运行情况下,每十分钟发送一次资源消耗情况。 |
|
作业状态变更 |
投递作业状态变更事件,包含状态变化前后信息。 |
作业状态实际发生变更时触发投递。 |
注意事项
-
投递至SLS前,需创建日志服务Project和Logstore,操作详情请参见使用LoongCollector采集并分析ECS文本日志。
-
投递至Kafka前,需在云消息队列Kafka版控制台创建实例和Topic,系统不会自动创建,详情请参见创建Kafka资源。
-
SLS投递配置与Kafka投递配置相互独立,可同时开启,也可仅开启其中一个,投递范围需在两个页签分别配置。
-
投递功能本身不收费,SLS的使用(日志库索引等)会产生流量自费,详情请参见计费概述。
-
SLS中查询、分析日志需要开启索引,开启后会产生索引流量和索引存储空间,费用说明请参见计费概述,是否创建索引由您自主决定。
-
您可设置日志库的服务器端加密方式,会话记录投递会继承该加密方式,具体操作请参见数据加密。
-
消息投递的范围包括启动日志、资源详情、运行事件全量投递、作业消耗资源和作业状态变更。作业运行日志的输出配置,请参见配置作业日志输出。
-
变更消息投递配置,会有最长30s的延迟。
操作步骤
步骤一:配置消息投递渠道
-
进入消息投递配置管理入口
-
登录实时计算管理控制台。
-
单击目标工作空间操作列下的控制台。
-
在左侧导航栏上,单击。
-
单击消息投递配置页签。
-
-
配置消息投递参数
根据目标投递渠道,完成对应页签的配置。
SLS投递配置
-
打开投递到日志服务SLS开关。
-
配置SLS相关信息。
配置项
说明
授权模式
-
STS Token:选择此模式,只能投递到Flink工作空间所在地域的SLS项目下的日志库。因此您仅需要再配置日志服务的项目和日志库。
-
AccessKey:选择此模式,支持投递到其他地域的SLS项目下的日志库。因此您需要再配置Region、Endpoint、AccessKeyId和AccessKeySecret。
日志服务项目
SLS项目名称。
日志服务日志库
SLS LogStore。
Region
SLS项目所在的地域。
当授权模式为AccessKey时,需要手动配置。各地域对应的Region和Endpoint地址,请参见服务接入点。
Endpoint
SLS服务的Endpoint地址。
当授权模式为STS Token时,系统会默认设置为Flink工作空间所在地域对应的endpoint;当授权模式为AccessKey时,需要您手动配置。
投递范围
具体消息内容,详情请参见下文字段含义说明。
AccessKeyId
阿里云账号的AccessKey ID和AccessKey Secret。
重要为了避免AK信息泄露,此处通过变量的方式来管理您的AK信息。您可以单击下拉箭头选择已创建的变量名,或单击参数右侧的
直接创建新变量。有关变量管理及如何查看AccessKey ID和AccessKey Secret值的详细信息,请参见变量管理和如何查看AccessKey ID和AccessKey Secret信息?
AccessKeySecret
-
-
单击保存。
Kafka投递配置
-
打开投递到消息队列Kafka开关。
-
配置Kafka相关信息。
配置项
说明
授权模式
-
STS Token:仅支持投递至与Flink工作空间同地域的Kafka实例,Kafka实例、Endpoint和Kafka Topic均支持下拉选择。
-
手动配置:需手动填写Endpoint和Kafka Topic。
Kafka实例
消息投递的目标Kafka实例。授权模式为STS Token时显示,可单击右侧的刷新图标重新拉取实例列表。
Endpoint
Kafka实例的接入点地址。授权模式为STS Token时,随所选实例自动填充;授权模式为手动配置时需手动填写。
Kafka Topic
消息投递的目标Topic,需预先创建,可单击Kafka控制台跳转创建。
授权模式为手动配置时,Topic名称长度为3~64个字符,仅包含英文、数字、半角句号(.)、短横线(-)和下划线(_),且至少包含1个英文或数字。
投递范围
选择需投递的消息类型,至少选择1项:启动日志、资源详情、运行事件全量投递、作业消耗资源、作业状态变更。
具体消息内容,详情请参见下文字段含义说明。
-
-
单击保存。
-
步骤二:查看投递消息
操作步骤
投递到SLS
-
单击日志服务项目右侧的打开SLS控制台。
-
查看原始日志的信息。
原始日志为JSON格式的资源配额(RESOURCE_QUOTA)信息,包含
namespaceTotalCpu(20.0)、namespaceTotalMemory(0.08 TiB)、namespaceUsedCpu(8.0)、namespaceUsedMemory(23.00 GiB)、resourceQueueName(default-queue)等字段,记录命名空间及资源队列的CPU与内存配额及使用量。
投递到Kafka
-
投递的消息以JSON格式写入步骤一配置的Kafka Topic。
-
登录云消息队列Kafka版控制台,在目标实例的消息查询页面查看投递的消息,操作详情请参见消息查询。也可自建Kafka客户端,使用配置的Endpoint消费该Topic。
字段含义说明
投递至SLS与Kafka的消息体结构一致,各消息类型的字段说明如下。
启动日志(JOB_START_LOG)
|
字段 |
含义 |
|
messageType |
消息类型,取值固定为JOB_START_LOG(作业启动日志)。 |
|
deploymentId |
已部署作业ID。 |
|
deploymentName |
已部署作业名称。 |
|
jobId |
作业实例ID。 |
|
tag |
作业标签(如无配置为空)。 |
|
length |
日志总长度。 |
|
offset |
用于日志分片时标记该条日志的起始位置。 |
|
content |
作业启动日志详情。 |
|
workspace |
工作空间ID。 |
|
namespace |
项目空间名称。 |
|
messageId |
消息ID。 |
|
timestamp |
时间戳。 |
资源详情(JOB_RESOURCE_QUOTA)
|
字段 |
含义 |
|
messageType |
消息类型,取值固定为RESOURCE_QUOTA(资源使用)。 |
|
namespaceTotalCpuMemory |
项目空间总共的资源CU数。 |
|
namespaceTotalCpu |
项目空间总共的资源CU数。 |
|
namespaceTotalMemory |
项目空间总共的内存资源量, |
|
namespaceUsedCpuMemory |
项目空间已使用的资源CU数。 |
|
namespaceUsedCpu |
项目空间已使用的资源CU数。 |
|
namespaceUsedMemory |
项目空间已使用的内存资源量。 |
|
resourceQueueName |
资源队列名称。 |
|
resourceQueueTotalCpuMemory |
资源队列总资源CU数。 |
|
resourceQueueTotalCpu |
资源队列总资源CU数。 |
|
resourceQueueTotalMemory |
资源队列总内存资源量。 |
|
resourceQueueUsedCpuMemory |
资源队列已使用的资源CU数。 |
|
resourceQueueUsedCpu |
资源队列已使用的资源CU数。 |
|
resourceQueueUsedMemory |
资源队列已使用的内存资源量。 |
|
workspace |
工作空间ID。 |
|
namespace |
项目空间名称。 |
|
messageId |
消息ID。 |
|
timestamp |
时间戳。 |
运行事件全量投递(JOB_EVENT)
|
字段 |
含义 |
|
messageType |
消息类型,取值固定为JOB_EVENT(运行事件)。 |
|
deploymentId |
已部署作业ID。 |
|
deploymentName |
已部署作业名称。 |
|
jobId |
作业实例ID。 |
|
tag |
作业标签(如无配置为空) |
|
eventId |
事件ID。 |
|
eventName |
事件名称。 |
|
content |
作业启动日志详情。 |
|
workspace |
工作空间ID。 |
|
namespace |
项目空间名称。 |
|
messageId |
消息ID。 |
|
timestamp |
时间戳。 |
作业消耗资源(JOB_RESOURCE_USAGE)
|
字段 |
含义 |
|
messageType |
消息类型,取值固定为JOB_RESOURCE_USAGE(作业消耗资源)。 |
|
deploymentId |
已部署作业ID。 |
|
deploymentName |
已部署作业名称。 |
|
jobId |
作业实例ID。 |
|
tag |
作业标签(如无配置为空)。 |
|
jobUsedCpu |
作业已使用CU数。 |
|
jobUsedMemory |
作业已使用内存数。 |
|
workspace |
工作空间ID。 |
|
namespace |
项目空间名称。 |
|
messageId |
消息ID。 |
|
timestamp |
时间戳。 |
作业状态变更(JOB_STATUS_CHANGED)
|
字段 |
含义 |
|
messageType |
消息类型,取值固定为JOB_STATUS_CHANGED(作业状态变更)。 |
|
messageId |
消息唯一标识。 |
|
timestamp |
消息产生的时间戳(毫秒)。 |
|
workspace |
工作空间名称。 |
|
namespace |
项目空间名称。 |
|
deploymentId |
已部署作业ID。 |
|
deploymentName |
已部署作业名称。 |
|
jobId |
作业实例ID。 |
|
previousStatus |
变更前的作业状态。 |
|
currentStatus |
变更后的作业状态。 |
|
tag |
作业标签(如无配置为空)。 |
|
additionalInfo |
附加信息。 |
相关文档
-
单作业的日志配置方法,请参见配置作业日志输出。
-
实时计算开发控制台上查看日志的操作方法详情,请参见查看启动和运行日志、查看运行事件、查看运行异常日志和查看历史作业实例日志。