[自研开源] MyData v0.8 数据集成之实时同步

2024-04-17 19:36

本文主要是介绍[自研开源] MyData v0.8 数据集成之实时同步,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

开源地址:gitee | github
详细介绍:MyData 基于 Web API 的数据集成平台

诚邀试用

可通过界面配置 实现API数据集成,减少集成工作,诚挚邀请更多用户试用;

系统的主要功能和优势:

  • 集成度高:通过API集成数据,没有数据库类型、开发技术等限制,有API即可对接;
  • 数据安全:无需开放业务系统数据库;
  • 零侵入性:无SDK,业务系统有API即可集成;
  • 可控性高:统一调度集成任务,对业务系统几乎无感知;
  • 复用性高:业务系统迁移、技术升级等,无需重复集成,API可用即可恢复集成;
  • 私有部署:支持本地部署,防止数据外泄;

特此承诺:

作为试用用户可享受 免费试用免费升级功能全程技术支持

前10位纳入实际项目使用的将成为永久免费用户

试用方式:联系微信,开通专属账号,开展数据集成;

截止日期:2024-08-30

wx

案例:电商场景 - 跨平台实时同步商品库存

接之前的定时同步案例,本次通过界面配置对接两个商城系统的webhook,实现商品库存的实时同步;

预期实现的流程如下:
在这里插入图片描述
主要过程:

  1. MyData提供两个接口,分别配置到两个系统的webhook;
  2. 当一方有库存变更时,调用webhook中配置的接口,接收业务数据;
  3. MyData接收并存储数据;
  4. 同时触发业务数据的消费任务,调用另一方接口推送更新的数据;

MyData配置过程

  1. 配置数据来源任务,接收trade webhook推送的数据

    • 创建提供数据类型的任务
    • 模式选择接收推送
    • 选择认证方式认证参数
    • 配置接口与业务数据的字段映射
      在这里插入图片描述
  2. 保存任务后,未任务生成了随机地址,点击复制
    在这里插入图片描述

  3. 在tradevine的webhook配置中 选择warehouse stock updated事件,填写目标地址和Header参数;
    在这里插入图片描述

  4. 然后配置数据消费任务,调用woo系统的接口 更新商品库存

    • 创建消费数据类型的任务
    • 模式选择调用API,并选择API接口
    • 选择订阅,并选择触发订阅的任务为上一个配置的任务
    • 单数据模式选择对象
    • 配置接口与业务数据的字段映射
      在这里插入图片描述
  5. 自此完成了一方的实时同步,另一方采用相同配置即可;

  6. 实时同步的执行过程

  • 当trade系统中 库存数量发生变化时,触发并调用webhook接口,以下是接口执行过程代码片段:
@Slf4j
@RestController
@AllArgsConstructor
@RequestMapping(MdConstant.API_PREFIX_MANAGE + "/integration")
public class IntegrationEndpoint {@Resourceprivate final ITaskService taskService;@Resourceprivate final JobExecutor jobExecutor;@PostMapping("/{task_url}")public R post(@PathVariable("task_url") String taskUrl, @RequestHeader HttpHeaders httpHeaders, @RequestBody String body) {log.info("integration post");return execute(taskUrl, httpHeaders, body);}@PutMapping("/{task_url}")public R put(@PathVariable("task_url") String taskUrl, @RequestHeader HttpHeaders httpHeaders, @RequestBody String body) {log.info("integration put");return execute(taskUrl, httpHeaders, body);}private R execute(String taskUrl, HttpHeaders httpHeaders, String body) {log.info("integration url : {}", taskUrl);log.info("integration headers : {}", httpHeaders);log.info("integration body : {}", body);Assert.notEmpty(taskUrl, "操作失败:地址 {} 无效!", taskUrl);Task task = taskService.findByApiUrl(taskUrl);Assert.notNull(task, "操作失败:地址 {} 无效!", taskUrl);Assert.equals(task.getTaskStatus(), MdConstant.TASK_STATUS_RUNNING, "操作失败:任务未启动 无法执行!");// 省略校验过程...// 执行任务流程 接收数据jobExecutor.acceptData(task, body);return R.success(StrUtil.format("任务 [{}] 开始执行,请查看日志。", task.getTaskName()));}
}
  • 任务流程中,使用accept到的json数据,根据字段映射进行解析和保存
// 将json按字段映射 解析为业务数据
jobDataService.parseProduceData(taskInfo, json);// 根据条件过滤数据
if (CollUtil.isNotEmpty(taskInfo.getDataFilters())) {taskInfo.appendLog("过滤业务数据开始");taskInfo.appendLog("过滤条件:{}", taskInfo.getDataFilters());jobDataFilterService.doFilter(taskInfo);taskInfo.appendLog("过滤后的剩余数据量:{}", taskInfo.getProduceDataList().size());if (CollUtil.isEmpty(taskInfo.getProduceDataList())) {taskInfo.appendLog("过滤后的没有业务数据,跳过后续处理");break;}
}// 保存业务数据
jobDataService.saveTaskData(taskInfo);// 更新环境变量
jobVarService.saveVarValue(taskInfo, json);// 递增分批参数
jobBatchService.incBatchParam(taskInfo);
  • 接收任务完成后,查询订阅该任务的消费任务,使用相同任务批次号 以便消费相同的业务数据
// 查询相同数据的订阅任务
List<Task> subTasks = taskService.listRunningSubTasks(taskInfo.getDataId(), taskInfo.getEnvId(), taskInfo.getId());
subTasks.forEach(task -> {TaskInfo subTaskInfo = build(task);// 订阅任务现在执行subTaskInfo.setStartTime(new Date());// 设置数据批次编号subTaskInfo.setDataBatchId(taskInfo.getDataBatchId());// 指定订阅任务,调用接口发送数据executeJob(subTaskInfo);
});
  1. 消费任务中,查询相同批次的数据,并调用API发送数据
// 订阅任务 使用任务批次号 查询数据
if (MdConstant.TASK_IS_SUBSCRIBED.equals(taskInfo.getIsSubscribed())) {// 构建数据批次查询条件 _MD_BATCH_ID_ = dataBatchIdBizDataFilter bizDataFilter = new BizDataFilter();bizDataFilter.setKey(MdConstant.DATA_COLUMN_BATCH_ID);bizDataFilter.setOp(MdConstant.DATA_OP_EQ);bizDataFilter.setValue(taskInfo.getDataBatchId());bizDataFilter.setType(MdConstant.TASK_FILTER_TYPE_VALUE);filters.add(bizDataFilter);
}
// 消费模式是调用API
if (MdConstant.TASK_CONSUME_MODE_API.equals(taskInfo.getConsumeMode())) {// 根据字段映射转换为api参数jobDataService.convertConsumeData(taskInfo);// 若消费任务是对象模式,则从字段映射中提取数据 替换url上的变量if (MdConstant.TASK_SINGLE_MODE_OBJECT.equals(taskInfo.getSingleMode())) {// 从url中解析出变量jobVarService.parseConsumeUrlVar(taskInfo);taskInfo.appendLog("替换API变量后 新地址为:url={}", taskInfo.getApiUrl());}// 调用api传输数据taskInfo.appendLog("调用API 获取数据,method={},url={},headers={},params={}", taskInfo.getApiMethod(), taskInfo.getApiUrl(), taskInfo.getReqHeaders(), taskInfo.getReqParams());String json = ApiUtil.write(taskInfo);// 更新环境变量jobVarService.saveVarValue(taskInfo, json);
}

这篇关于[自研开源] MyData v0.8 数据集成之实时同步的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/912657

相关文章

【服务器运维】MySQL数据存储至数据盘

查看磁盘及分区 [root@MySQL tmp]# fdisk -lDisk /dev/sda: 21.5 GB, 21474836480 bytes255 heads, 63 sectors/track, 2610 cylindersUnits = cylinders of 16065 * 512 = 8225280 bytesSector size (logical/physical)

SpringBoot集成Netty,Handler中@Autowired注解为空

最近建了个技术交流群,然后好多小伙伴都问关于Netty的问题,尤其今天的问题最特殊,功能大概是要在Netty接收消息时把数据写入数据库,那个小伙伴用的是 Spring Boot + MyBatis + Netty,所以就碰到了Handler中@Autowired注解为空的问题 参考了一些大神的博文,Spring Boot非controller使用@Autowired注解注入为null的问题,得到

vue项目集成CanvasEditor实现Word在线编辑器

CanvasEditor实现Word在线编辑器 官网文档:https://hufe.club/canvas-editor-docs/guide/schema.html 源码地址:https://github.com/Hufe921/canvas-editor 前提声明: 由于CanvasEditor目前不支持vue、react 等框架开箱即用版,所以需要我们去Git下载源码,拿到其中两个主

SQL Server中,查询数据库中有多少个表,以及数据库其余类型数据统计查询

sqlserver查询数据库中有多少个表 sql server 数表:select count(1) from sysobjects where xtype='U'数视图:select count(1) from sysobjects where xtype='V'数存储过程select count(1) from sysobjects where xtype='P' SE

时间服务器中,适用于国内的 NTP 服务器地址,可用于时间同步或 Android 加速 GPS 定位

NTP 是什么?   NTP 是网络时间协议(Network Time Protocol),它用来同步网络设备【如计算机、手机】的时间的协议。 NTP 实现什么目的?   目的很简单,就是为了提供准确时间。因为我们的手表、设备等,经常会时间跑着跑着就有误差,或快或慢的少几秒,时间长了甚至误差过分钟。 NTP 服务器列表 最常见、熟知的就是 www.pool.ntp.org/zo

探索Elastic Search:强大的开源搜索引擎,详解及使用

🎬 鸽芷咕:个人主页  🔥 个人专栏: 《C++干货基地》《粉丝福利》 ⛺️生活的理想,就是为了理想的生活! 引入 全文搜索属于最常见的需求,开源的 Elasticsearch (以下简称 Elastic)是目前全文搜索引擎的首选,相信大家多多少少的都听说过它。它可以快速地储存、搜索和分析海量数据。就连维基百科、Stack Overflow、

数据时代的数字企业

1.写在前面 讨论数据治理在数字企业中的影响和必要性,并介绍数据治理的核心内容和实践方法。作者强调了数据质量、数据安全、数据隐私和数据合规等方面是数据治理的核心内容,并介绍了具体的实践措施和案例分析。企业需要重视这些方面以实现数字化转型和业务增长。 数字化转型行业小伙伴可以加入我的星球,初衷成为各位数字化转型参考库,星球内容每周更新 个人工作经验资料全部放在这里,包含数据治理、数据要

如何在Java中处理JSON数据?

如何在Java中处理JSON数据? 大家好,我是免费搭建查券返利机器人省钱赚佣金就用微赚淘客系统3.0的小编,也是冬天不穿秋裤,天冷也要风度的程序猿!今天我们将探讨在Java中如何处理JSON数据。JSON(JavaScript Object Notation)作为一种轻量级的数据交换格式,在现代应用程序中被广泛使用。Java通过多种库和API提供了处理JSON的能力,我们将深入了解其用法和最佳

两个基因相关性CPTAC蛋白组数据

目录 蛋白数据下载 ①蛋白数据下载 1,TCGA-选择泛癌数据  2,TCGA-TCPA 3,CPTAC(非TCGA) ②蛋白相关性分析 1,数据整理 2,蛋白相关性分析 PCAS在线分析 蛋白数据下载 CPTAC蛋白组学数据库介绍及数据下载分析 – 王进的个人网站 (jingege.wang) ①蛋白数据下载 可以下载泛癌蛋白数据:UCSC Xena (xena

BD错误集锦8——在集成Spring MVC + MyBtis编写mapper文件时需要注意格式 You have an error in your SQL syntax

报错的文件 <?xml version="1.0" encoding="UTF-8" ?><!DOCTYPE mapperPUBLIC "-//mybatis.org//DTD Mapper 3.0//EN""http://mybatis.org/dtd/mybatis-3-mapper.dtd"><mapper namespace="com.yuan.dao.YuanUserDao"><!