feat(sync): 增加pulsar消息最后接收时间

用于判断pulsar有多久没有传入消息了
This commit is contained in:
v-zhangjc9
2024-05-09 09:03:00 +08:00
parent 103bde5cdc
commit 0bf3d17009
6 changed files with 101 additions and 11 deletions

View File

@@ -45,11 +45,11 @@ public interface SQLConstants {
*/
String ALIAS_A = _alias_.getAlias() + "." + ALIAS_O;
/**
* 字段 flink_job_id 原始值 flink_job_id Flink Job Id
* 字段 flink_job_id 原始值 flink_job_id Flink Job Id
*/
String FLINK_JOB_ID_O = "flink_job_id";
/**
* 字段 flink_job_id 别名值 tacti.flink_job_id Flink Job Id
* 字段 flink_job_id 别名值 tacti.flink_job_id Flink Job Id
*/
String FLINK_JOB_ID_A = _alias_.getAlias() + "." + FLINK_JOB_ID_O;
/**
@@ -618,6 +618,22 @@ public interface SQLConstants {
* 字段 type 别名值 tahcm.type 类型
*/
String TYPE_A = _alias_.getAlias() + "." + TYPE_O;
/**
* 字段 source_schema 原始值 source_schema 库名
*/
String SOURCE_SCHEMA_O = "source_schema";
/**
* 字段 source_schema 别名值 tahcm.source_schema 库名
*/
String SOURCE_SCHEMA_A = _alias_.getAlias() + "." + SOURCE_SCHEMA_O;
/**
* 字段 source_table 原始值 source_table 表名
*/
String SOURCE_TABLE_O = "source_table";
/**
* 字段 source_table 别名值 tahcm.source_table 表名
*/
String SOURCE_TABLE_A = _alias_.getAlias() + "." + SOURCE_TABLE_O;
/**
* 字段 compaction_plan_instant 原始值 compaction_plan_instant 压缩计划时间点
*/
@@ -910,6 +926,14 @@ public interface SQLConstants {
* 字段 source_start_time 别名值 tahss.source_start_time 同步启动时间
*/
String SOURCE_START_TIME_A = _alias_.getAlias() + "." + SOURCE_START_TIME_O;
/**
* 字段 source_receive_time 原始值 source_receive_time pulsar最后接收消息的时间
*/
String SOURCE_RECEIVE_TIME_O = "source_receive_time";
/**
* 字段 source_receive_time 别名值 tahss.source_receive_time pulsar最后接收消息的时间
*/
String SOURCE_RECEIVE_TIME_A = _alias_.getAlias() + "." + SOURCE_RECEIVE_TIME_O;
/**
* 字段 source_checkpoint_time 原始值 source_checkpoint_time 同步检查时间
*/