feat(info-query): 增加更新pulsar backlog数据的接口
This commit is contained in:
@@ -44,4 +44,13 @@ public class SyncStateController {
|
||||
public void saveCompactionId(@RequestParam("id") String id) {
|
||||
syncStateService.saveSyncState(id);
|
||||
}
|
||||
|
||||
@GetMapping("/sync_state/save_pulsar_back_log")
|
||||
public void savePulsarBacklog(
|
||||
@RequestParam("flink_job_id") Long flinkJobId,
|
||||
@RequestParam("alias") String alias,
|
||||
@RequestParam("back_log") Long backlog
|
||||
) {
|
||||
syncStateService.savePulsarBacklog(flinkJobId, alias, backlog);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,7 +17,9 @@ import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.retry.annotation.Retryable;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.*;
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.TbAppCollectTableInfo;
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.TbAppFlinkJobConfig;
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.TbAppHudiSyncState;
|
||||
|
||||
/**
|
||||
* Sync State
|
||||
@@ -117,4 +119,21 @@ public class SyncStateService extends BaseService {
|
||||
"-1:-1:-1"
|
||||
);
|
||||
}
|
||||
|
||||
public void savePulsarBacklog(Long flinkJobId, String alias, Long backlog) {
|
||||
mysqlJdbcTemplate.update(
|
||||
SqlBuilder
|
||||
.insertInto(
|
||||
TbAppHudiSyncState._origin_,
|
||||
TbAppHudiSyncState.ID_O,
|
||||
TbAppHudiSyncState.PULSAR_BACK_LOG_O
|
||||
)
|
||||
.values()
|
||||
.addValue(Q, Q)
|
||||
.onDuplicateKeyUpdateColumn(TbAppHudiSyncState.PULSAR_BACK_LOG_O)
|
||||
.precompileSql(),
|
||||
StrUtil.format("{}-{}", flinkJobId, alias),
|
||||
backlog
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,7 +9,12 @@ import cn.hutool.core.util.StrUtil;
|
||||
import cn.hutool.db.sql.SqlUtil;
|
||||
import com.lanyuanxiaoyao.service.common.SQLConstants;
|
||||
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.*;
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.TbAppCollectTableInfo;
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.TbAppFlinkJobConfig;
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.TbAppGlobalConfig;
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.TbAppHudiJobConfig;
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.TbAppHudiSyncState;
|
||||
import static com.lanyuanxiaoyao.service.common.SQLConstants.HudiCollectBuild.TbAppYarnJobConfig;
|
||||
|
||||
/**
|
||||
* @author lanyuanxiaoyao
|
||||
@@ -64,17 +69,10 @@ public class SqlBuilderTests {
|
||||
|
||||
public static void main(String[] args) {
|
||||
System.out.println(SqlUtil.formatSql(
|
||||
generateTableMetaList(
|
||||
SqlBuilder.select(
|
||||
TbAppFlinkJobConfig.ID_A,
|
||||
TbAppCollectTableInfo.ALIAS_A,
|
||||
TbAppCollectTableInfo.TGT_HDFS_PATH_A,
|
||||
TbAppCollectTableInfo.SRC_PULSAR_ADDR_A,
|
||||
TbAppCollectTableInfo.SRC_TOPIC_A
|
||||
),
|
||||
null,
|
||||
null
|
||||
).build()
|
||||
SqlBuilder.update(TbAppHudiSyncState._origin_)
|
||||
.set(TbAppHudiSyncState.PULSAR_BACK_LOG_O, "?")
|
||||
.whereEq(TbAppHudiSyncState.ID_O, Column.as("?"))
|
||||
.build()
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user