首页 文章 精选 留言 我的

精选列表

搜索[精准预约],共7977篇文章
优秀的个人博客,低调大师

风控开发指南:Java集成人车关系核验实现精准合规审查

在现代智慧物流与同城货运平台的运力调度场景中,确认司机真实身份与接单车辆的精确对应关系,是确保运营合规与货物安全的核心前提。过去,平台通常依赖人工审核司机上传的行驶证和驾驶证照片。这种传统核查模式耗时费力,且难以动态排查低频非活跃账号或非存续异常主体,极易导致实际承运人与注册资质不符的合规隐患。

优秀的个人博客,低调大师

风控开发指南:Python集成车辆静态信息查询实现精准合规审查

在物流供应链金融及智能调度的核心数据平台中,通常会使用 Python 生态下的分布式任务流来管理商用车队的入驻评估。在评估阶段,核对车辆排量、燃料属性与真实生产方等信息,常常需借助线下补充资料与人工登记。这不仅拉长了入驻周期,还极易混入非存续异常主体或导致资质信息错配,带来不确定的授信及履约隐患。

优秀的个人博客,低调大师

风控开发指南:PHP集成车辆静态信息查询实现精准合规审查

在二手车经销商CRM后台管理系统及数字化车辆交易服务中,对拟收购车辆真实参数、出厂配置与合规排放标准进行前置准入校验,是控制履约风险与保障资产交易安全的核心环节。传统的人工核查机动车登记证书(绿本)照片或依赖线下经验进行判定的方式流程繁琐、效率较低,且难以实时防范信息不匹配资质带来的隐患,容易在车辆估值模型和跨区域流通中造成严重的数据阻塞。

优秀的个人博客,低调大师

风控开发指南:Java集成车辆静态信息查询实现精准合规审查

在大型企业车队资产管理与物流供应链微服务架构中,准确掌握车队的车辆底盘级静态信息,对于车队资源分配、环保合规与全生命周期管理至关重要。传统的车队资产入库环节,往往依赖运维人员人工登记纸质出厂证明或行驶证。不仅录入过程繁琐、容易产生数据孤岛,还极易混入非存续异常主体或异动企业提供的虚假档案材料,从而在后续的道路运输和尾气环保检测中埋下巨大的履约隐患。企业迫切需要引入标准化的数据通道以替代人工校验,构建自动化的云原生数字资产系统。

优秀的个人博客,低调大师

​Cursor 升级 Tab 模型,实时强化学习提升开发者建议精准

Cursor 是一款基于人工智能的编程平台,最近宣布对其 Tab 模型进行了升级。Tab 模型是为开发者提供自动补全建议的系统。此次升级显著减少了低质量建议的数量,提高了建议的准确性。具体来说,新的 Tab 模型相比于之前的版本,建议数量减少了21%,而接受率提高了28%。 Cursor 在其博客中表示,实现高接受率不仅仅是让模型变得更智能,还需要懂得何时提供建议、何时不提供。为了应对这一挑战,Cursor 考虑了训练一个单独的模型,用于预测某个建议是否会被接受。该公司引用了一项2022年的研究,指出这种方法在 GitHub Copilot 中取得了成功。研究中采用了逻辑回归过滤器,分析编程语言、最近的接受历史和训练字符等特征,将那些得分较低的建议隐藏起来。 然而,Cursor 认为这种解决方案虽然可以预测用户接受建议的概率,但希望有一个更通用的机制,能够重用 Tab 模型学到的强大代码表示。Cursor 希望通过改变 Tab 模型的结构,避免在最初就产生低质量建议,而不是后续再进行过滤。 因此,Cursor 采用了策略梯度方法,这是一种强化学习的方法。当用户接受建议时,模型会得到奖励;当建议被拒绝时,模型会受到惩罚;而在选择保持沉默时则不会得到任何反馈。此方法需要 “在线” 数据,即从当前使用的模型收集的反馈。Cursor 通过每天多次向用户部署新的检查点,并迅速基于新交互对模型进行再训练,来解决这一问题。 Cursor 表示,当前从部署检查点到收集数据的过程仅需1.5到2小时,这在 AI 行业中已经算是较快,但仍有进一步加速的空间。该公司的 Tab 模型每天处理超过4亿个请求,Cursor 希望这一改进能够提升开发者的编码体验,并计划在未来进一步开发这些方法。 在线强化学习是该领域最令人兴奋的方向之一,一位在 OpenAI 从事后训练的工程师在社交媒体上对此表示赞赏,称 Cursor 似乎是第一个成功在大规模上实施该技术的公司。 不久前,Cursor 的母公司 Anysphere 宣布融资9亿美元,估值达99亿美元,并推出了一项月费200美元的 “超值” 计划,承诺提供20倍于20美元月费 “专业版” 的使用量。此外,Cursor 还在同月进行了平台更新,新增了自动代码审查、记忆功能和一键设置模型上下文协议服务器的功能。

优秀的个人博客,低调大师

SeaTunnel 如何给 MySQL 表做“精准切片”?一篇读懂 CDC 分片黑科技

作者 | 梁尧博(本文由AI辅助生成) 概述 Apache SeaTunnel MySQL CDC连接器为了实现并行读取,需要将大表切分成多个分片(Split)。对于非主键表,连接器提供了多种智能切分策略来保证数据的完整性和读取效率。本系列将详细介绍 Apache SeaTunnel 支持的核心切分策略,切分策略机制及实现方式,并对比各个切分策略的优劣势。 1. 切分列选择策略 1.1 选择优先级 1. 用户配置的snapshotSplitColumn(建议是唯一键) 2. 主键列(按数据类型优先级选择) 3. 唯一键列(按数据类型优先级选择) 4. 无可用列 → 单分片策略 1.2 支持的数据类型 MySQL CDC连接器支持的切分列类型: 根据AbstractJdbcSourceChunkSplitter.isEvenlySplitColumn()方法的实现: // AbstractJdbcSourceChunkSplitter.isEvenlySplitColumn() switch (fromDbzColumn(splitColumn).getSqlType()) { case TINYINT: case SMALLINT: case INT: case BIGINT: case DECIMAL: case STRING: return true; default: return false; } ✅ 支持的类型: 数值类型:TINYINT、SMALLINT、INT、BIGINT、DECIMAL 字符串类型:STRING(使用哈希切分) ❌ 注意:MySQL CDC不支持日期时间类型切分 不支持的类型: ❌ DATE:不支持作为切分列 ❌ DATETIME:不支持作为切分列 ❌ TIMESTAMP:不支持作为切分列 ❌ TIME:不支持作为切分列 对比:普通JDBC连接器的支持情况 值得注意的是,普通JDBC连接器(DynamicChunkSplitter)支持DATE类型切分: // DynamicChunkSplitter支持的类型 switch (splitColumnType.getSqlType()) { case TINYINT: case SMALLINT: case INT: case BIGINT: case DECIMAL: case DOUBLE: case FLOAT: return evenlyColumnSplitChunks(table, splitColumnName, min, max, chunkSize); case STRING: // 字符串切分逻辑 case DATE: // ✅ 普通JDBC支持DATE类型 return dateColumnSplitChunks(table, splitColumnName, min, max, chunkSize); } 实际影响和解决方案 如果表只有日期时间类型的索引列,MySQL CDC将: 无法找到合适的切分列 回退到单分片模式 失去并行读取的优势 1.3 数据类型优先级 // 优先级:1最高,6最低 TINYINT(1) > SMALLINT(2) > INT(3) > BIGINT(4) > DECIMAL(5) > STRING(6) 2. 切分策略决策机制 SeaTunnel通过一套精密的决策算法来确定使用哪种切分策略,这个决策过程基于数据分布特征和表大小等因素。 2.1 决策流程概述 // 核心决策逻辑位于 AbstractJdbcSourceChunkSplitter.generateSplits() public Collection<SnapshotSplit> generateSplits(JdbcConnection jdbc, TableId tableId) { // 1. 获取配置参数 final int chunkSize = sourceConfig.getSplitSize(); // 默认:8096 final double distributionFactorUpper = sourceConfig.getDistributionFactorUpper(); // 默认:100.0 final double distributionFactorLower = sourceConfig.getDistributionFactorLower(); // 默认:0.05 final int sampleShardingThreshold = sourceConfig.getSampleShardingThreshold(); // 默认:1000 // 2. 检查切分列类型 if (isEvenlySplitColumn(splitColumn)) { // 3. 查询近似行数和计算分布因子 long approximateRowCnt = queryApproximateRowCnt(jdbc, tableId); double distributionFactor = calculateDistributionFactor(tableId, min, max, approximateRowCnt); // 4. 判断数据分布是否均匀 boolean dataIsEvenlyDistributed = distributionFactor >= distributionFactorLower && distributionFactor <= distributionFactorUpper; if (dataIsEvenlyDistributed) { // 均匀切分策略 return splitEvenlySizedChunks(...); } else { // 5. 检查是否需要采样策略 int shardCount = (int) (approximateRowCnt / chunkSize); if (sampleShardingThreshold < shardCount) { // 采样切分策略 return efficientShardingThroughSampling(...); } else { // 不均匀切分策略 return splitUnevenlySizedChunks(...); } } } else { // 字符串类型:不均匀切分策略 return splitUnevenlySizedChunks(...); } } 2.2 分布因子计算 核心公式: distributionFactor = (MAX - MIN + 1) / approximateRowCount 计算逻辑: protected double calculateDistributionFactor(TableId tableId, Object min, Object max, long approximateRowCnt) { if (approximateRowCnt == 0) { return Double.MAX_VALUE; // 空表处理 } BigDecimal difference = ObjectUtils.minus(max, min); final BigDecimal subRowCnt = difference.add(BigDecimal.valueOf(1)); double distributionFactor = subRowCnt.divide( new BigDecimal(approximateRowCnt), 4, ROUND_CEILING).doubleValue(); return distributionFactor; } 分布因子含义: factor ≈ 1.0:数据分布理想,ID连续且无空隙 factor > 100:数据稀疏,ID范围远大于行数(如:ID 1-100万,但只有1000行) factor < 0.05:数据密集,多行共享相似的ID值(如:时间戳列,同一秒内多条记录) 2.3 决策条件详解 条件1:切分列类型检查 // 支持均匀切分的数据类型 private boolean isEvenlySplitColumn(Column splitColumn) { return splitColumn.isNumeric() || splitColumn.isTemporalType(); } 条件2:数据分布均匀性判断 boolean dataIsEvenlyDistributed = doubleCompare(distributionFactor, distributionFactorLower) >= 0 && doubleCompare(distributionFactor, distributionFactorUpper) <= 0; // 即:0.05 ≤ distributionFactor ≤ 100 条件3:采样策略触发条件 int shardCount = (int) (approximateRowCnt / chunkSize); if (sampleShardingThreshold < shardCount) { // 预估分片数超过1000时,启用采样策略 } 2.4 实际决策示例 示例1:理想均匀分布 表:user_orders 切分列:order_id (BIGINT) 数据范围:1 - 100,000 行数:100,000 chunkSize:10,000 计算: distributionFactor = (100000 - 1 + 1) / 100000 = 1.0 判断:0.05 ≤ 1.0 ≤ 100 → 数据均匀分布 结果:使用均匀切分策略,生成10个分片 示例2:数据稀疏,触发采样 表:big_transactions 切分列:transaction_id (BIGINT) 数据范围:1 - 10,000,000 行数:50,000 chunkSize:1,000 计算: distributionFactor = (10000000 - 1 + 1) / 50000 = 200 预估分片数 = 50000 / 1000 = 50 判断:200 > 100 → 数据分布不均匀 50 < 1000 → 不触发采样 结果:使用不均匀切分策略 示例3:大表触发采样策略 表:log_events 切分列:event_id (BIGINT) 数据范围:1 - 1,00,000 行数:5,000,000 chunkSize:1,000 计算: distributionFactor = (100000 - 1 + 1) / 5000000 = 0.2 预估分片数 = 5000000 / 1000 = 5000 判断:0.02 < 0.05 → 数据分布不均匀 5000 > 1000 → 触发采样策略 结果:使用采样切分策略,采样后生成5000个分片 示例4:时间戳列密集分布(假设支持时间戳类型) 表:sensor_data 切分列:timestamp (TIMESTAMP) 数据范围:2023-01-01 00:00:00 - 2023-01-01 01:00:00 (3600秒) 行数:1,000,000 chunkSize:10,000 计算: distributionFactor = 3600 / 1000000 = 0.0036 判断:0.0036 < 0.05 → 数据分布不均匀 预估分片数 = 1000000 / 10000 = 100 < 1000 结果:使用不均匀切分策略 2.5 策略选择总结 条件组合 分布因子范围 预估分片数 选择策略 适用场景 数值列 + 均匀分布 [0.05, 100] 任意 均匀切分 自增ID,均匀分布的数值 数值列 + 不均匀 + 小表 <0.05 或 >100 ≤1000 不均匀切分 稀疏ID,时间戳密集 数值列 + 不均匀 + 大表 <0.05 或 >100 >1000 采样切分 大表且分布极不均匀 字符串列 不适用 任意 不均匀切分 字符串类型切分列 3. 三种核心切分策略 3.1 均匀切分(Evenly Sized Chunks) 适用场景: 数据分布均匀的数值列 判断条件: // 分布因子计算 distributionFactor = (max - min + 1) / approximateRowCount // 均匀分布判断 distributionFactorLower <= distributionFactor <= distributionFactorUpper // 默认:0.05 <= distributionFactor <= 100 切分逻辑: // 动态chunk大小计算 dynamicChunkSize = Math.max((int)(distributionFactor * chunkSize), 1) // 切分范围计算 chunkStart = null chunkEnd = min + dynamicChunkSize while (chunkEnd <= max) { splits.add(ChunkRange.of(chunkStart, chunkEnd)) chunkStart = chunkEnd chunkEnd = chunkEnd + dynamicChunkSize } // 添加最后一个分片 splits.add(ChunkRange.of(chunkStart, null)) 示例: 表:user_table,主键:id,范围:1-10000,行数:10000 distributionFactor = (10000-1+1)/10000 = 1.0 chunkSize = 1000,dynamicChunkSize = 1000 分片结果: Split1: [null, 1000] // id <= 1000 Split2: [1000, 2000] // 1000 < id <= 2000 Split3: [2000, 3000] // 2000 < id <= 3000 ... Split10: [9000, null] // id > 9000 3.2 不均匀切分(Unevenly Sized Chunks) 适用场景: 数据分布不均匀或非数值列 切分逻辑: // 连续查询下一个chunk的最大值 Object chunkStart = null Object chunkEnd = queryNextChunkMax(jdbc, min, tableId, splitColumn, max, chunkSize) while (chunkEnd != null && chunkEnd <= max) { splits.add(ChunkRange.of(chunkStart, chunkEnd)) chunkStart = chunkEnd chunkEnd = queryNextChunkMax(jdbc, chunkEnd, tableId, splitColumn, max, chunkSize) } splits.add(ChunkRange.of(chunkStart, null)) SQL示例: -- 查询下一个chunk的最大值 SELECT MAX(split_column) FROM ( SELECT split_column FROM table_name WHERE split_column >= ? ORDER BY split_column LIMIT ? ) t 示例: 表:order_table,切分列:create_time,chunkSize=1000 查询过程: 1. 查询前1000行的最大create_time → '2023-01-15 10:30:00' 2. 查询接下来1000行的最大create_time → '2023-02-20 15:45:00' 3. 继续查询... 分片结果: Split1: [null, '2023-01-15 10:30:00'] Split2: ['2023-01-15 10:30:00', '2023-02-20 15:45:00'] Split3: ['2023-02-20 15:45:00', '2023-03-25 09:20:00'] ... 3.3 采样切分(Sampling-based Sharding) 适用场景: 大表且数据分布极不均匀 触发条件: // 当预估分片数超过阈值时启用 int shardCount = (int)(approximateRowCount / chunkSize) if (sampleShardingThreshold < shardCount) { // 使用采样切分 } // 默认阈值:1000 采样逻辑: // 采样数据 Object[] sampleData = sampleDataFromColumn(jdbc, tableId, splitColumn, inverseSamplingRate) // 计算每个分片的样本数 double approxSamplePerShard = (double)sampleData.length / shardCount // 根据样本数据确定分片边界 for (int i = 0; i < shardCount; i++) { Object chunkStart = lastEnd Object chunkEnd = (i < shardCount - 1) ? sampleData[(int)((i + 1) * approxSamplePerShard)] : null splits.add(ChunkRange.of(chunkStart, chunkEnd)) } 示例: 表:big_table,行数:1000万,chunkSize=10000,预估分片数=1000 inverseSamplingRate=1000(采样率1/1000) 采样过程: 1. 从表中采样10000行数据 2. 将采样数据按切分列排序 3. 根据1000个分片需求,每10个样本确定一个分片边界 分片结果:基于采样数据的分布确定边界 4. SQL查询案例详解与代码实现 4.1 核心SQL查询方法映射 SQL查询类型 对应方法 实现类 具体作用 MIN/MAX查询 queryMinMax() MySqlUtils.java 获取切分列的最小最大值 行数统计查询 queryApproximateRowCnt() MySqlUtils.java 获取表的近似行数 动态边界查询 queryNextChunkMax() MySqlUtils.java 不均匀切分的边界计算 采样数据查询 sampleDataFromColumn() MySqlUtils.java 采样策略的数据采集 字符串哈希查询 hashModForField() MysqlDialect.java 字符串类型的哈希切分 4.2 均匀切分策略SQL示例与实现 4.2.1 数值类型(BIGINT)- 自增ID表 表结构: CREATE TABLE user_orders ( order_id BIGINT AUTO_INCREMENT PRIMARY KEY, user_id INT, amount DECIMAL(10,2), created_at TIMESTAMP ); -- 数据范围:order_id 1-100000,共100000行 切分过程SQL: -- 1. 查询MIN/MAX值 SELECT MIN(order_id), MAX(order_id) FROM user_orders; -- 结果:MIN=1, MAX=100000 -- 2. 查询近似行数(MySQL) SHOW TABLE STATUS LIKE 'user_orders'; -- 结果:Rows=100000 -- 3. 计算分布因子:(100000-1+1)/100000 = 1.0 -- 判断:0.05 ≤ 1.0 ≤ 100 → 均匀切分 -- 4. 生成分片SQL(chunkSize=10000) -- Split 0: SELECT * FROM user_orders WHERE order_id >= 1 AND order_id < 10001; -- Split 1: SELECT * FROM user_orders WHERE order_id >= 10001 AND order_id < 20001; -- Split 2: SELECT * FROM user_orders WHERE order_id >= 20001 AND order_id < 30001; -- ...继续到Split 9 -- Split 9: SELECT * FROM user_orders WHERE order_id >= 90001 AND order_id <= 100000; 对应代码实现: public static Object[] queryMinMax(JdbcConnection jdbc, TableId tableId, String columnName) throws SQLException { final String minMaxQuery = String.format( "SELECT MIN(%s), MAX(%s) FROM %s", quote(columnName), quote(columnName), quote(tableId)); return jdbc.queryAndMap( minMaxQuery, rs -> { if (!rs.next()) { throw new SQLException( String.format( "No result returned after running query [%s]", minMaxQuery)); } return rowToArray(rs, 2); }); } public static long queryApproximateRowCnt(JdbcConnection jdbc, TableId tableId) throws SQLException { final String useDatabaseStatement = String.format("USE %s;", quote(tableId.catalog())); final String rowCountQuery = String.format("SHOW TABLE STATUS LIKE '%s';", tableId.table()); jdbc.execute(useDatabaseStatement); return jdbc.queryAndMap( rowCountQuery, rs -> { if (!rs.next() || rs.getMetaData().getColumnCount() < 5) { throw new SQLException( String.format( "No result returned after running query [%s]", rowCountQuery)); } return rs.getLong(5); // Rows列 }); } 4.2.2 普通JDBC连接器的DATE类型切分(仅供参考) // 普通JDBC连接器支持DATE类型切分,CDC模式不支持 switch (splitColumnType.getSqlType()) { case DATE: return dateColumnSplitChunks(table, splitColumnName, min, max, chunkSize); } 表结构: CREATE TABLE daily_reports ( id BIGINT AUTO_INCREMENT PRIMARY KEY, -- ✅ MySQL CDC使用此列切分 report_date DATE, -- ❌ MySQL CDC不支持此列切分 revenue DECIMAL(15,2), created_at TIMESTAMP ); MySQL CDC的实际切分方式: -- MySQL CDC会使用id列进行切分,而不是report_date -- 1. 查询MIN/MAX值(基于id列) SELECT MIN(id), MAX(id) FROM daily_reports; -- 结果:MIN=1, MAX=365 -- 2. 查询近似行数 SHOW TABLE STATUS LIKE 'daily_reports'; -- 结果:Rows=365 -- 3. 计算分布因子(基于id列) -- 分布因子:(365-1+1)/365 = 1 -- 判断:1 < 100 → 使用均匀切分 -- 4. 生成分片SQL(基于id列) -- Split 0: SELECT * FROM daily_reports WHERE id >= 1 AND id < 92; -- Split 1: SELECT * FROM daily_reports WHERE id >= 92 AND id < 183; -- Split 2: SELECT * FROM daily_reports WHERE id >= 183 AND id < 274; -- Split 3: SELECT * FROM daily_reports WHERE id >= 274 AND id <= 365; 4.3 不均匀切分策略SQL示例与实现 4.3.1 数值类型 - 稀疏ID表 表结构: CREATE TABLE sparse_transactions ( transaction_id BIGINT PRIMARY KEY, account_id INT, amount DECIMAL(15,2) ); -- 数据特点:ID范围1-1000000,但只有1000行数据 切分过程SQL: -- 1. 查询MIN/MAX值 SELECT MIN(transaction_id), MAX(transaction_id) FROM sparse_transactions; -- 结果:MIN=1, MAX=1000000 -- 2. 查询近似行数 SHOW TABLE STATUS LIKE 'sparse_transactions'; -- 结果:Rows=1000 -- 3. 计算分布因子:(1000000-1+1)/1000 = 1000 -- 判断:1000 > 100 → 数据分布不均匀 -- 4. 不均匀切分SQL生成(chunkSize=100) -- 动态查询每个分片的边界值 -- 查询第1个分片的结束值 SELECT MAX(transaction_id) FROM ( SELECT transaction_id FROM sparse_transactions WHERE transaction_id >= 1 ORDER BY transaction_id ASC LIMIT 100 ) AS T; -- 假设结果:15000 -- Split 0: SELECT * FROM sparse_transactions WHERE transaction_id >= 1 AND transaction_id < 15000; -- 查询第2个分片的结束值 SELECT MAX(transaction_id) FROM ( SELECT transaction_id FROM sparse_transactions WHERE transaction_id >= 15000 ORDER BY transaction_id ASC LIMIT 100 ) AS T; -- 假设结果:35000 -- Split 1: SELECT * FROM sparse_transactions WHERE transaction_id >= 15000 AND transaction_id < 35000; -- 继续此过程直到所有数据被分片 对应代码实现: public static Object queryNextChunkMax( JdbcConnection jdbc, TableId tableId, String splitColumnName, int chunkSize, Object includedLowerBound) throws SQLException { String quotedColumn = quote(splitColumnName); String query = String.format( "SELECT MAX(%s) FROM (" + "SELECT %s FROM %s WHERE %s >= ? ORDER BY %s ASC LIMIT %s" + ") AS T", quotedColumn, quotedColumn, quote(tableId), quotedColumn, quotedColumn, chunkSize); return jdbc.prepareQueryAndMap( query, ps -> ps.setObject(1, includedLowerBound), rs -> { if (!rs.next()) { throw new SQLException( String.format( "No result returned after running query [%s]", query)); } return rs.getObject(1); }); } 4.3.2 字符串类型切分策略详解 ⚠️ 重要说明: MySQL CDC连接器仅支持哈希取模方式切分字符串,而普通JDBC连接器还支持字符集编码方式。 4.3.2.1 MySQL CDC的字符串切分(哈希取模) 表结构: CREATE TABLE customer_profiles ( customer_code VARCHAR(50) PRIMARY KEY, name VARCHAR(100), email VARCHAR(100) ); -- 数据特点:customer_code为字符串,如 'CUST001', 'CUST002'等 切分过程SQL: -- 字符串类型使用哈希取模切分 -- 1. 查询总行数 SHOW TABLE STATUS LIKE 'customer_profiles'; -- 结果:假设40000行 -- 2. 计算分片数:40000 / 8096 ≈ 5个分片 -- 3. 生成哈希取模分片SQL(使用MD5哈希,SeaTunnel默认) -- Split 0: SELECT * FROM customer_profiles WHERE ABS(MD5(customer_code) % 5) = 0; -- Split 1: SELECT * FROM customer_profiles WHERE ABS(MD5(customer_code) % 5) = 1; -- Split 2: SELECT * FROM customer_profiles WHERE ABS(MD5(customer_code) % 5) = 2; -- Split 3: SELECT * FROM customer_profiles WHERE ABS(MD5(customer_code) % 5) = 3; -- Split 4: SELECT * FROM customer_profiles WHERE ABS(MD5(customer_code) % 5) = 4; 4.3.2.2 普通JDBC的字符串切分(字符集编码) ⚠️ 注意: 以下内容仅适用于普通JDBC连接器,MySQL CDC不支持此方式。 配置参数: # 启用字符集编码切分模式(仅普通JDBC) split.string_split_mode = charset_based split.string_split_mode_collate = "0123456789abcdefghijklmnopqrstuvwxyz" public enum StringSplitMode { SAMPLE("sample"), // 默认:采样方式(实际还是哈希) CHARSET_BASED("charset_based"); // 字符集编码方式 } 字符集编码算法核心: // 字符串编码为BigInteger的核心算法 public static BigInteger encodeStringToNumericRange( String str, int maxLength, boolean paddingAtEnd, boolean isCaseInsensitive, String orderedCharset, int radix) { // 1. 字符串转换为ASCII索引表示 String asciiString = stringToAsciiString(str, maxLength, paddingAtEnd, isCaseInsensitive, orderedCharset); // 2. 解析为基数数组 int[] baseArray = parseBaseNumber(asciiString); // 3. 转换为BigInteger(类似进制转换) return toDecimal(baseArray, radix); } 字符集编码切分示例: -- 假设字符集:"0123456789abcdefghijklmnopqrstuvwxyz"(37个字符) -- 字符串"cust001"编码过程: -- 1. 'c'→12, 'u'→30, 's'→28, 't'→29, '0'→0, '0'→0, '1'→1 -- 2. BigInteger = 12*37^6 + 30*37^5 + 28*37^4 + 29*37^3 + 0*37^2 + 0*37^1 + 1*37^0 -- 3. 按数值范围进行均匀切分 -- Split 0: 编码值范围 [0, 50000000000) SELECT * FROM customer_profiles WHERE encode_string_to_bigint(customer_code) >= 0 AND encode_string_to_bigint(customer_code) < 50000000000; -- Split 1: 编码值范围 [50000000000, 100000000000) SELECT * FROM customer_profiles WHERE encode_string_to_bigint(customer_code) >= 50000000000 AND encode_string_to_bigint(customer_code) < 100000000000; 4.3.2.3 两种字符串切分方式对比 特性 哈希取模切分 字符集编码切分 支持连接器 ✅ MySQL CDC<br/>✅ 普通JDBC ❌ MySQL CDC<br/>✅ 普通JDBC 切分原理 哈希函数+取模 字符串→数值编码→范围切分 数据分布 随机分布(哈希特性) 按字典序分布 性能特点 计算简单,速度快 编码复杂,但支持范围查询 适用场景 数据随机分布需求 需要保持字符串顺序的场景 配置复杂度 简单(无需配置) 复杂(需配置字符集) 实际应用建议: MySQL CDC场景:只能使用哈希取模,无其他选择 普通JDBC场景: 默认使用哈希取模(性能更好) 需要保持字符串顺序时使用字符集编码 字符集编码适合固定格式的字符串(如编号、代码等) 对应代码实现: default String hashModForField(String fieldName, int mod) { return "ABS(MD5(" + quoteIdentifier(fieldName) + ") % " + mod + ")"; } 不同数据库的哈希实现: @Override public String hashModForField(String nativeType, String fieldName, int mod) { String quoteFieldName = quoteIdentifier(fieldName); if (StringUtils.isNotBlank(nativeType)) { quoteFieldName = convertType(quoteFieldName, nativeType); } return "(ABS(HASHTEXT(" + quoteFieldName + ")) % " + mod + ")"; } 4.4 采样切分策略SQL示例与实现 4.4.1 超大表采样切分 表结构: CREATE TABLE big_log_events ( event_id BIGINT, user_id INT, event_type VARCHAR(50), timestamp TIMESTAMP, INDEX idx_event_id (event_id) ); -- 数据特点:5000万行,event_id分布不均匀 切分过程SQL: -- 1. 查询MIN/MAX值 SELECT MIN(event_id), MAX(event_id) FROM big_log_events; -- 结果:MIN=1, MAX=1000000 -- 2. 查询近似行数 SHOW TABLE STATUS LIKE 'big_log_events'; -- 结果:Rows=50000000 -- 3. 计算分布因子:(1000000-1+1)/50000000 = 0.02 -- 预估分片数:50000000/8096 ≈ 6177 -- 判断:0.02 < 0.05 且 6177 > 1000 → 触发采样切分 -- 4. 采样查询(采样率 1/1000,即inverseSamplingRate=1000) SELECT event_id FROM big_log_events WHERE MOD((event_id - (SELECT MIN(event_id) FROM big_log_events)), 1000) = 0 ORDER BY event_id; -- 采样约50000行数据 -- 5. 根据采样结果计算分片边界 -- 假设采样后得到边界值:[1, 15000, 28000, 45000, 67000, ...] -- 6. 生成最终分片SQL -- Split 0: SELECT * FROM big_log_events WHERE event_id >= 1 AND event_id < 15000; -- Split 1: SELECT * FROM big_log_events WHERE event_id >= 15000 AND event_id < 28000; -- Split 2: SELECT * FROM big_log_events WHERE event_id >= 28000 AND event_id < 45000; -- 继续直到所有分片 对应代码实现: public static Object[] sampleDataFromColumn( JdbcConnection jdbc, TableId tableId, String columnName, int inverseSamplingRate) throws SQLException { final String minQuery = String.format( "SELECT %s FROM %s WHERE MOD((%s - (SELECT MIN(%s) FROM %s)), %s) = 0 ORDER BY %s", quote(columnName), quote(tableId), quote(columnName), quote(columnName), quote(tableId), inverseSamplingRate, quote(columnName)); return jdbc.queryAndMap( minQuery, resultSet -> { List<Object> results = new ArrayList<>(); while (resultSet.next()) { results.add(resultSet.getObject(1)); } return results.toArray(); }); } 5. SQL查询模式总结 5.1 各策略SQL查询模式对比 策略类型 数据类型 SQL查询模式 示例 均匀切分 数值类型 WHERE col >= start AND col < end WHERE order_id >= 1 AND order_id < 10001 均匀切分 字符串 哈希取模查询 WHERE ABS(CRC32(name) % 4) = 0 不均匀切分 数值类型 动态边界查询 SELECT MAX(id) FROM (SELECT id FROM table WHERE id >= ? ORDER BY id LIMIT ?) 不均匀切分 字符串 哈希取模查询 WHERE ABS(MD5(name) % 4) = 0 采样切分 数值类型 采样+边界查询 WHERE MOD((id - (SELECT MIN(id) FROM table)), 1000) = 0 采样切分 字符串 字符串采样查询 WHERE ABS(CRC32(name) % 1000) = 0 5.2 多数据库SQL差异对比 数据库 哈希函数 行数统计 分页语法 MySQL MD5(field) SHOW TABLE STATUS LIMIT n PostgreSQL HASHTEXT(field) pg_class.reltuples LIMIT n SQL Server HASHBYTES('MD5', field) sys.dm_db_partition_stats TOP n Oracle ORA_HASH(field) all_tables.num_rows ROWNUM <= n 6. 性能优化与配置 6.1 分布因子调优 # 分布因子配置 chunk-key.even-distribution.factor.upper-bound = 100.0 # 上限,默认100.0 chunk-key.even-distribution.factor.lower-bound = 0.05 # 下限,默认0.05 参数说明: chunk-key.even-distribution.factor.upper-bound:均匀分布因子上限,用于判断数据是否均匀分布 chunk-key.even-distribution.factor.lower-bound:均匀分布因子下限,计算公式:(MAX(id) - MIN(id) + 1) / row count 6.2 采样策略调优 # 采样配置 sample-sharding.threshold = 1000 # 采样阈值,默认1000 inverse-sampling.rate = 1000 # 采样率倒数,默认1000 参数说明: sample-sharding.threshold:触发采样切分策略的预估分片数阈值 inverse-sampling.rate:采样率的倒数,例如1000表示1/1000的采样率 6.3 快照切分配置 # 快照切分配置 snapshot.split.size = 8096 # 分片大小,默认8096行 snapshot.fetch.size = 1024 # 每次拉取大小,默认1024行 参数说明: snapshot.split.size:表快照的分片大小(行数) snapshot.fetch.size:读取表快照时每次轮询的最大拉取大小 6.4 核心配置参数 参数名 类型 默认值 说明 snapshot.split.size Integer 8096 表快照的分片大小(行数) snapshot.fetch.size Integer 1024 读取快照时每次拉取的最大行数 chunk-key.even-distribution.factor.upper-bound Double 100.0 均匀分布因子上限 chunk-key.even-distribution.factor.lower-bound Double 0.05 均匀分布因子下限 sample-sharding.threshold Integer 1000 采样切分阈值 inverse-sampling.rate Integer 1000 采样率倒数 server-id String 随机生成 数据库客户端的唯一ID server-time-zone String UTC 数据库服务器的会话时区 connect.timeout.ms Duration 30000 连接超时时间(毫秒) connect.max-retries Integer 3 最大重试次数 connection.pool.size Integer 20 JDBC连接池大小 6.5 配置示例 source { Mysql-CDC { # 基础连接配置 url = "jdbc:mysql://localhost:3306/test" username = "root" password = "123456" table-names = ["test.user_table"] # 快照切分配置 snapshot.split.size = 8096 snapshot.fetch.size = 1024 # 分布因子配置 chunk-key.even-distribution.factor.upper-bound = 100.0 chunk-key.even-distribution.factor.lower-bound = 0.05 # 采样策略配置 sample-sharding.threshold = 1000 inverse-sampling.rate = 1000 # 连接配置 server-id = "5400" server-time-zone = "Asia/Shanghai" connect.timeout.ms = 30000 connect.max-retries = 3 connection.pool.size = 20 # 启动模式 startup.mode = "initial" # 其他配置 exactly_once = false format = "DEFAULT" } } 7. 切分策略控制与现场应用 7.1 策略控制参数总结 通过调整关键参数,可以精确控制SeaTunnel使用哪种切分策略,以应对不同的现场场景: 1. 强制使用均匀切分策略 适用场景: 数据分布相对均匀,追求最佳并行性能 source { Mysql-CDC { url = "jdbc:mysql://localhost:3306/test" username = "root" password = "123456" table-names = ["test.uniform_table"] # 强制均匀切分配置 chunk-key.even-distribution.factor.upper-bound = 10000.0 # 大幅提高上限 chunk-key.even-distribution.factor.lower-bound = 0.001 # 大幅降低下限 sample-sharding.threshold = 100000 # 极高阈值避免采样 snapshot.split.size = 8096 # 标准分片大小 } } 2. 强制使用不均匀切分策略 适用场景: 数据分布不均,但表不是特别大 source { Mysql-CDC { url = "jdbc:mysql://localhost:3306/test" username = "root" password = "123456" table-names = ["test.sparse_table"] # 强制不均匀切分配置 chunk-key.even-distribution.factor.upper-bound = 0.1 # 极低上限 chunk-key.even-distribution.factor.lower-bound = 0.1 # 极低下限 sample-sharding.threshold = 100000 # 极高阈值避免采样 snapshot.split.size = 5000 # 适中分片大小 } } 3. 强制使用采样切分策略 适用场景: 超大表,需要高效切分 source { Mysql-CDC { url = "jdbc:mysql://localhost:3306/test" username = "root" password = "123456" table-names = ["test.huge_table"] # 强制采样切分配置 chunk-key.even-distribution.factor.upper-bound = 0.01 # 极低上限 chunk-key.even-distribution.factor.lower-bound = 0.01 # 极低下限 sample-sharding.threshold = 100 # 极低阈值强制采样 inverse-sampling.rate = 500 # 提高采样率 snapshot.split.size = 10000 # 较大分片大小 } } 4. 避免采样策略(业务库压力大) 适用场景: 大表但业务库压力大,不能使用采样 source { Mysql-CDC { url = "jdbc:mysql://localhost:3306/test" username = "root" password = "123456" table-names = ["test.large_table"] # 避免采样配置 chunk-key.even-distribution.factor.upper-bound = 1000.0 # 放宽上限 chunk-key.even-distribution.factor.lower-bound = 0.001 # 放宽下限 sample-sharding.threshold = 50000 # 极高阈值 snapshot.split.size = 50000 # 大分片减少总数 connection.pool.size = 5 # 减少连接数 snapshot.fetch.size = 1024 # 控制拉取大小 } } 5. 高并行性能优化 适用场景: 追求最大并行度和处理速度 source { Mysql-CDC { url = "jdbc:mysql://localhost:3306/test" username = "root" password = "123456" table-names = ["test.performance_table"] # 高并行配置 snapshot.split.size = 2000 # 小分片增加并行度 snapshot.fetch.size = 2048 # 增加拉取大小 connection.pool.size = 30 # 增加连接池 chunk-key.even-distribution.factor.upper-bound = 1000.0 # 优先均匀切分 sample-sharding.threshold = 10000 # 适中阈值 } } 7.2 参数决策矩阵 场景类型 分片大小 上限因子 下限因子 采样阈值 策略结果 数据均匀,追求性能 2000-8096 10000.0 0.001 100000 均匀切分 数据稀疏,中等表 5000-10000 0.1 0.1 100000 不均匀切分 超大表,允许采样 10000+ 0.01 0.01 100 采样切分 大表,业务库压力大 50000+ 1000.0 0.001 50000 避免采样 高并行需求 2000 1000.0 0.001 10000 均匀切分 7.3 核心控制参数说明 策略选择控制: chunk-key.even-distribution.factor.upper-bound: 控制是否使用均匀切分 chunk-key.even-distribution.factor.lower-bound: 控制分布判断的敏感度 sample-sharding.threshold: 控制是否触发采样策略 性能调优控制: snapshot.split.size: 控制并行度和内存使用 snapshot.fetch.size: 控制数据库查询压力 connection.pool.size: 控制数据库连接压力 inverse-sampling.rate: 控制采样精度 8. 总结 SeaTunnel MySQL CDC连接器的表切分机制通过以下核心组件实现: AbstractJdbcSourceChunkSplitter:核心切分逻辑 MySqlUtils:MySQL特定的SQL查询实现 JdbcDialect:数据库方言支持 三种切分策略: 均匀切分:适用于数据分布均匀的数值和日期类型 不均匀切分:适用于数据分布稀疏的场景 采样切分:适用于超大表的高效切分 决策机制: 通过分布因子判断数据分布特征 根据表大小选择合适的切分策略 支持多种数据类型的特殊处理 通过精确的参数控制,这套机制能够应对各种复杂的现场场景,确保MySQL CDC在处理各种规模和类型的表时都能实现高效、均衡的数据切分。 附 | 切分流程 SeaTunnel切分策略决策参数

优秀的个人博客,低调大师

技术解析|Doris Connector 结合 Flink CDC 实现 MySQL 分库分表 Exactly Once精准接入

1. 概述 在实际业务系统中为了解决单表数据量大带来的各种问题,我们通常采用分库分表的方式对库表进行拆分,以达到提高系统的吞吐量。 但是这样给后面数据分析带来了麻烦,这个时候我们通常试将业务数据库的分库分表同步到数据仓库时,将这些分库分表的数据,合并成一个库,一个表。便于我们后面的数据分析 本篇文档我们就演示怎么基于Flink CDC 并结合 Apache Doris Flink Connector 及 Doris Stream Load的两阶段提交,实现MySQL数据库分库分表实时高效的接入到 Apache Doris 数据仓库中进行分析。 1.1 什么是CDC CDC是(Change Data Capture 变更数据获取)的简称。 核心思想是,监测并捕获数据库的变动(包括数据 或 数据表的插入INSERT、更新UPDATE、删除DELETE等),将这些变更按发生的顺序完整记录下来,写入到消息中间件中以供其他服务进行订阅及消费。 CDC 技术应用场景也非常广泛,包括: ● 数据分发,将一个数据源分发给多个下游,常用于业务解耦、微服务。 ● 数据集成,将分散异构的数据源集成到数据仓库中,消除数据孤岛,便于后续的分析。 ● 数据迁移,常用于数据库备份、容灾等。 1.2 为什么选择Flink CDC Flink CDC 基于数据库日志的Change Data Caputre 技术,实现了全量和增量的一体化读取能力,并借助 Flink 优秀的管道能力和丰富的上下游生态,支持捕获多种数据库的变更,并将这些变更实时同步到下游存储。 目前,Flink CDC 的上游已经支持了 MySQL、MariaDB、PG、Oracle、MongoDB 、Oceanbase、TiDB、SQLServer等数据库。 Flink CDC 的下游则更加丰富,支持写入 Kafka、Pulsar 消息队列,也支持写入 Hudi、Iceberg 、Doris等,支持写入各种数据仓库及数据湖中。 同时,通过 Flink SQL 原生支持的 Changelog 机制,可以让 CDC 数据的加工变得非常简单。用户可以通过 SQL 便能实现数据库全量和增量数据的清洗、打宽、聚合等操作,极大地降低了用户门槛。 此外, Flink DataStream API 支持用户编写代码实现自定义逻辑,给用户提供了深度定制业务的自由度 Flink CDC 技术的核心是支持将表中的全量数据和增量数据做实时一致性的同步与加工,让用户可以方便地获每张表的实时一致性快照。比如一张表中有历史的全量业务数据,也有增量的业务数据在源源不断写入,更新。Flink CDC 会实时抓取增量的更新记录,实时提供与数据库中一致性的快照,如果是更新记录,会更新已有数据。如果是插入记录,则会追加到已有数据,整个过程中,Flink CDC 提供了一致性保障,即不重不丢。 FLink CDC 如下优势: Flink 的算子和 SQL 模块更为成熟和易用 Flink 作业可以通过调整算子并行度的方式,轻松扩展处理能力 Flink 支持高级的状态后端(State Backends),允许存取海量的状态数据 Flink 提供更多的 Source 和 Sink 等生态支持 Flink 有更大的用户基数和活跃的支持社群,问题更容易解决 而且 Flink Table / SQL 模块将数据库表和变动记录流(例如 CDC 的数据流)看做是同一事物的两面,因此内部提供的 Upsert 消息结构(+I 表示新增、-U 表示记录更新前的值、+U 表示记录更新后的值,-D 表示删除)可以与 Debezium 等生成的变动记录一一对应。 1.3 什么是Apache Doris Apache Doris是一个现代化的MPP分析型数据库产品。仅需亚秒级响应时间即可获得查询结果,有效地支持实时数据分析。Apache Doris的分布式架构非常简洁,易于运维,并且可以支持10PB以上的超大数据集。 Apache Doris可以满足多种数据分析需求,例如固定历史报表,实时数据分析,交互式数据分析和探索式数据分析等。令您的数据分析工作更加简单高效! 1.4 Two-phase commit 1.4.1 什么是 two-phase commit (2PC) 在分布式系统中,为了让每个节点都能够感知到其他节点的事务执行状况,需要引入一个中心节点来统一处理所有节点的执行逻辑,这个中心节点叫做协调者(coordinator),被中心节点调度的其他业务节点叫做参与者(participant)。 2PC将分布式事务分成了两个阶段,两个阶段分别为提交请求(投票)和提交(执行)。协调者根据参与者的响应来决定是否需要真正地执行事务,具体流程如下。 提交请求(投票)阶段 协调者向所有参与者发送prepare请求与事务内容,询问是否可以准备事务提交,并等待参与者的响应。 参与者执行事务中包含的操作,并记录undo日志(用于回滚)和redo日志(用于重放),但不真正提交。 参与者向协调者返回事务操作的执行结果,执行成功返回yes,否则返回no。 提交(执行)阶段 分为成功与失败两种情况。 若所有参与者都返回yes,说明事务可以提交: 协调者向所有参与者发送commit请求。 参与者收到commit请求后,将事务真正地提交上去,并释放占用的事务资源,并向协调者返回ack。 协调者收到所有参与者的ack消息,事务成功完成。 若有参与者返回no或者超时未返回,说明事务中断,需要回滚: 协调者向所有参与者发送rollback请求。 参与者收到rollback请求后,根据undo日志回滚到事务执行前的状态,释放占用的事务资源,并向协调者返回ack。 协调者收到所有参与者的ack消息,事务回滚完成。 1.4 Flink 2PC Flink作为流式处理引擎,自然也提供了对exactly once语义的保证。端到端的exactly once语义,是输入、处理逻辑、输出三部分协同作用的结果。Flink内部依托检查点机制和轻量级分布式快照算法ABS保证exactly once。而要实现精确一次的输出逻辑,则需要施加以下两种限制之一:幂等性写入(idempotent write)、事务性写入(transactional write)。 预提交阶段的流程 每当需要做checkpoint时,JobManager就在数据流中打入一个屏障(barrier),作为检查点的界限。屏障随着算子链向下游传递,每到达一个算子都会触发将状态快照写入状态后端的动作。当屏障到达Kafka sink后,通过KafkaProducer.flush()方法刷写消息数据,但还未真正提交。接下来还是需要通过检查点来触发提交阶段 提交阶段流程 只有在所有检查点都成功完成这个前提下,写入才会成功。这符合前文所述2PC的流程,其中JobManager为协调者,各个算子为参与者(不过只有sink一个参与者会执行提交)。一旦有检查点失败,notifyCheckpointComplete()方法就不会执行。如果重试也不成功的话,最终会调用abort()方法回滚事务 1.5 Doris Stream Load 2PC 1.5.1 Stream load Stream load 是Apache Doris 提供的一个同步的导入方式,用户通过发送 HTTP 协议发送请求将本地文件或数据流导入到 Doris 中。Stream load 同步执行导入并返回导入结果。用户可直接通过请求的返回体判断本次导入是否成功。 Stream load 主要适用于导入本地文件,或通过程序导入数据流中的数据。 使用方法,用户通过Http Client 进行操作,也可以使用Curl命令进行 curl --location-trusted -u user:passwd [-H ""...] -T data.file -H "label:label" -XPUT http://fe_host:http_port/api/{db}/{table}/_stream_load 这里为了是防止用户重复导入相同的数据,使用了导入任务标识label。强烈推荐用户同一批次数据使用相同的 label。这样同一批次数据的重复请求只会被接受一次,保证了 At-Most-Once 1.5.2 Stream load 2PC Aapche Doris 最早的Stream Load 是没有两阶段提交的,导入数据的时候直接通过 Stream Load 的 http 接口完成数据导入,只有成功和失败。 这种在正常情况下是没有问题的,在分布式环境下可能为因为某一个导入任务是失败导致两端数据不一致的情况,特别是在Doris Flink Connector里,之前的Doris Flink Connector 数据导入失败需要用户自己控制,做异常处理,比如如果导入失败之后,将数据保存到指定的地方(例如Kafka),然后人工手动处理。 如果Flink Job因为其他为题突然挂掉,这样会造成部分数据成功,部分数据失败,而且失败的数据因为没有checkpoint,重新启动Job也没办法重新消费失败的数据,造成两端数据不一致 为了解决上面的这些问题,保证两端数据一致性,我们实现了Doris Stream Load 2PC,原理如下: 提交分成两个阶段 第一阶段,提交数据写入任务,这个时候数据写入成功后,数据状态是不可见的,事务状态是PRECOMMITTED 数据写入成功之后,用户触发Commit操作,将事务状态变成VISIBLE,这个时候数据可以查询到 如果用户要方式这一批数据只需要通过事务ID,对事务触发abort操作,这批数据将会被自动删除掉 1.5.3 Stream load 2PC使用方式 在be.conf中配置disable_stream_load_2pc=false(重启生效) 并且 在 HEADER 中声明 two_phase_commit=true 。 发起预提交: curl --location-trusted -u user:passwd -H "two_phase_commit:true" -T test.txt http://fe_host:http_port/api/{db}/{table}/_stream_load 触发事务Commit操作 curl -X PUT --location-trusted -u user:passwd -H "txn_id:18036" -H "txn_operation:commit" http://fe_host:http_port/api/{db}/_stream_load_2pc 对事物触发abort操作 curl -X PUT --location-trusted -u user:passwd -H "txn_id:18037" -H "txn_operation:abort" http://fe_host:http_port/api/{db}/_stream_load_2pc 1.6 Doris Flink Connector 2PC 我们之前提供了Doris Flink Connector ,支持对Doris表数据的读,Upsert、delete(Unique key模型),但是存在可能因为Job失败或者其他异常情况导致两端数据不一致的问题。 为了解决这些问题,我们基于FLink 2PC 和Doris Stream Load 2PC对Doris Connector进行了改造升级,保证两端exactly once。 我们会在内存中维护读写的buffer,在启动的时候,开启写入,并异步的提交,期间通过http chunked的方式持续的将数据写入到BE,直到Checkpoint的时候,停止写入,这样做的好处是避免用户频繁提交http带来的开销,Checkpoint完成后会开启下一阶段的写入 在这个Checkpoint期间,可能是多个task任务同时在写一张表的数据,这些我们都会在这个Checkpoint期间对应一个全局的label,在checkpoint的时候将这个label对应的写入数据的事务进行统一的一次提交,将数据状态变成可见, 如果失败 Flink 在重启的时候会对这些数据通过checkpoint进行回放。 这样就可以保证Doris两端数据的一致 2. 系统架构 下面我们通过一个完整示例来看怎么去通过Doris Flink Connector最新版本(支持两阶段提交),来完成整合Flink CDC实现MySQL分库分表实时采集入库 这里我们通过Flink CDC 来完成MySQL分库分表数据采集 然后通过Doris Flink Connector来完成数据的入库 最后利用Doris的高并发、高性能的OLAP分析计算能力对外提供数据服务 3. MySQL 安装配置 3.1 安装MySQL 快速使用Docker安装配置Mysql,具体参照下面的连接 https://segmentfault.com/a/1190000021523570 3.2 开启Mysql binlog 进入 Docker 容器修改/etc/my.cnf 文件,在 [mysqld] 下面添加以下内容, log_bin=mysql_bin binlog-format=Row server-id=1 然后重启Mysql systemctl restart mysqld 3.3 准备数据 这里演示我们准备了两个库emp_1,emp_2,每个库下面各种主备了两张表employees_1,employees_2。并给出了一下初始化数据 CREATE DATABASE emp_1; USE emp_1; CREATE TABLE employees_1 ( emp_no INT NOT NULL, birth_date DATE NOT NULL, first_name VARCHAR(14) NOT NULL, last_name VARCHAR(16) NOT NULL, gender ENUM ('M','F') NOT NULL, hire_date DATE NOT NULL, PRIMARY KEY (emp_no) ); ​ INSERT INTO `employees_1` VALUES (10001,'1953-09-02','Georgi','Facello','M','1986-06-26'), (10002,'1964-06-02','Bezalel','Simmel','F','1985-11-21'), (10003,'1959-12-03','Parto','Bamford','M','1986-08-28'), (10004,'1954-05-01','Chirstian','Koblick','M','1986-12-01'), (10005,'1955-01-21','Kyoichi','Maliniak','M','1989-09-12'), (10006,'1953-04-20','Anneke','Preusig','F','1989-06-02'), (10007,'1957-05-23','Tzvetan','Zielinski','F','1989-02-10'), (10008,'1958-02-19','Saniya','Kalloufi','M','1994-09-15'), (10009,'1952-04-19','Sumant','Peac','F','1985-02-18'), (10010,'1963-06-01','Duangkaew','Piveteau','F','1989-08-24'), (10011,'1953-11-07','Mary','Sluis','F','1990-01-22'), (10012,'1960-10-04','Patricio','Bridgland','M','1992-12-18'), (10013,'1963-06-07','Eberhardt','Terkki','M','1985-10-20'), (10014,'1956-02-12','Berni','Genin','M','1987-03-11'), (10015,'1959-08-19','Guoxiang','Nooteboom','M','1987-07-02'), (10016,'1961-05-02','Kazuhito','Cappelletti','M','1995-01-27'), (10017,'1958-07-06','Cristinel','Bouloucos','F','1993-08-03'), (10018,'1954-06-19','Kazuhide','Peha','F','1987-04-03'), (10019,'1953-01-23','Lillian','Haddadi','M','1999-04-30'), (10020,'1952-12-24','Mayuko','Warwick','M','1991-01-26'), (10021,'1960-02-20','Ramzi','Erde','M','1988-02-10'), (10022,'1952-07-08','Shahaf','Famili','M','1995-08-22'), (10023,'1953-09-29','Bojan','Montemayor','F','1989-12-17'), (10024,'1958-09-05','Suzette','Pettey','F','1997-05-19'), (10025,'1958-10-31','Prasadram','Heyers','M','1987-08-17'), (10026,'1953-04-03','Yongqiao','Berztiss','M','1995-03-20'), (10027,'1962-07-10','Divier','Reistad','F','1989-07-07'), (10028,'1963-11-26','Domenick','Tempesti','M','1991-10-22'), (10029,'1956-12-13','Otmar','Herbst','M','1985-11-20'), (10030,'1958-07-14','Elvis','Demeyer','M','1994-02-17'), (10031,'1959-01-27','Karsten','Joslin','M','1991-09-01'), (10032,'1960-08-09','Jeong','Reistad','F','1990-06-20'), (10033,'1956-11-14','Arif','Merlo','M','1987-03-18'), (10034,'1962-12-29','Bader','Swan','M','1988-09-21'), (10035,'1953-02-08','Alain','Chappelet','M','1988-09-05'), (10036,'1959-08-10','Adamantios','Portugali','M','1992-01-03'); ​ CREATE TABLE employees_2 ( emp_no INT NOT NULL, birth_date DATE NOT NULL, first_name VARCHAR(14) NOT NULL, last_name VARCHAR(16) NOT NULL, gender ENUM ('M','F') NOT NULL, hire_date DATE NOT NULL, PRIMARY KEY (emp_no) ); ​ INSERT INTO `employees_2` VALUES (10037,'1963-07-22','Pradeep','Makrucki','M','1990-12-05'), (10038,'1960-07-20','Huan','Lortz','M','1989-09-20'), (10039,'1959-10-01','Alejandro','Brender','M','1988-01-19'), (10040,'1959-09-13','Weiyi','Meriste','F','1993-02-14'), (10041,'1959-08-27','Uri','Lenart','F','1989-11-12'), (10042,'1956-02-26','Magy','Stamatiou','F','1993-03-21'), (10043,'1960-09-19','Yishay','Tzvieli','M','1990-10-20'), (10044,'1961-09-21','Mingsen','Casley','F','1994-05-21'), (10045,'1957-08-14','Moss','Shanbhogue','M','1989-09-02'), (10046,'1960-07-23','Lucien','Rosenbaum','M','1992-06-20'), (10047,'1952-06-29','Zvonko','Nyanchama','M','1989-03-31'), (10048,'1963-07-11','Florian','Syrotiuk','M','1985-02-24'), (10049,'1961-04-24','Basil','Tramer','F','1992-05-04'), (10050,'1958-05-21','Yinghua','Dredge','M','1990-12-25'), (10051,'1953-07-28','Hidefumi','Caine','M','1992-10-15'), (10052,'1961-02-26','Heping','Nitsch','M','1988-05-21'), (10053,'1954-09-13','Sanjiv','Zschoche','F','1986-02-04'), (10054,'1957-04-04','Mayumi','Schueller','M','1995-03-13'); ​ ​ CREATE DATABASE emp_2; ​ USE emp_2; ​ CREATE TABLE employees_1 ( emp_no INT NOT NULL, birth_date DATE NOT NULL, first_name VARCHAR(14) NOT NULL, last_name VARCHAR(16) NOT NULL, gender ENUM ('M','F') NOT NULL, hire_date DATE NOT NULL, PRIMARY KEY (emp_no) ); ​ ​ INSERT INTO `employees_1` VALUES (10055,'1956-06-06','Georgy','Dredge','M','1992-04-27'), (10056,'1961-09-01','Brendon','Bernini','F','1990-02-01'), (10057,'1954-05-30','Ebbe','Callaway','F','1992-01-15'), (10058,'1954-10-01','Berhard','McFarlin','M','1987-04-13'), (10059,'1953-09-19','Alejandro','McAlpine','F','1991-06-26'), (10060,'1961-10-15','Breannda','Billingsley','M','1987-11-02'), (10061,'1962-10-19','Tse','Herber','M','1985-09-17'), (10062,'1961-11-02','Anoosh','Peyn','M','1991-08-30'), (10063,'1952-08-06','Gino','Leonhardt','F','1989-04-08'), (10064,'1959-04-07','Udi','Jansch','M','1985-11-20'), (10065,'1963-04-14','Satosi','Awdeh','M','1988-05-18'), (10066,'1952-11-13','Kwee','Schusler','M','1986-02-26'), (10067,'1953-01-07','Claudi','Stavenow','M','1987-03-04'), (10068,'1962-11-26','Charlene','Brattka','M','1987-08-07'), (10069,'1960-09-06','Margareta','Bierman','F','1989-11-05'), (10070,'1955-08-20','Reuven','Garigliano','M','1985-10-14'), (10071,'1958-01-21','Hisao','Lipner','M','1987-10-01'), (10072,'1952-05-15','Hironoby','Sidou','F','1988-07-21'), (10073,'1954-02-23','Shir','McClurg','M','1991-12-01'), (10074,'1955-08-28','Mokhtar','Bernatsky','F','1990-08-13'), (10075,'1960-03-09','Gao','Dolinsky','F','1987-03-19'), (10076,'1952-06-13','Erez','Ritzmann','F','1985-07-09'), (10077,'1964-04-18','Mona','Azuma','M','1990-03-02'), (10078,'1959-12-25','Danel','Mondadori','F','1987-05-26'), (10079,'1961-10-05','Kshitij','Gils','F','1986-03-27'), (10080,'1957-12-03','Premal','Baek','M','1985-11-19'), (10081,'1960-12-17','Zhongwei','Rosen','M','1986-10-30'), (10082,'1963-09-09','Parviz','Lortz','M','1990-01-03'), (10083,'1959-07-23','Vishv','Zockler','M','1987-03-31'), (10084,'1960-05-25','Tuval','Kalloufi','M','1995-12-15'); ​ ​ CREATE TABLE employees_2( emp_no INT NOT NULL, birth_date DATE NOT NULL, first_name VARCHAR(14) NOT NULL, last_name VARCHAR(16) NOT NULL, gender ENUM ('M','F') NOT NULL, hire_date DATE NOT NULL, PRIMARY KEY (emp_no) ); ​ INSERT INTO `employees_2` VALUES (10085,'1962-11-07','Kenroku','Malabarba','M','1994-04-09'), (10086,'1962-11-19','Somnath','Foote','M','1990-02-16'), (10087,'1959-07-23','Xinglin','Eugenio','F','1986-09-08'), (10088,'1954-02-25','Jungsoon','Syrzycki','F','1988-09-02'), (10089,'1963-03-21','Sudharsan','Flasterstein','F','1986-08-12'), (10090,'1961-05-30','Kendra','Hofting','M','1986-03-14'), (10091,'1955-10-04','Amabile','Gomatam','M','1992-11-18'), (10092,'1964-10-18','Valdiodio','Niizuma','F','1989-09-22'), (10093,'1964-06-11','Sailaja','Desikan','M','1996-11-05'), (10094,'1957-05-25','Arumugam','Ossenbruggen','F','1987-04-18'), (10095,'1965-01-03','Hilari','Morton','M','1986-07-15'), (10096,'1954-09-16','Jayson','Mandell','M','1990-01-14'), (10097,'1952-02-27','Remzi','Waschkowski','M','1990-09-15'), (10098,'1961-09-23','Sreekrishna','Servieres','F','1985-05-13'), (10099,'1956-05-25','Valter','Sullins','F','1988-10-18'), (10100,'1953-04-21','Hironobu','Haraldson','F','1987-09-21'), (10101,'1952-04-15','Perla','Heyers','F','1992-12-28'), (10102,'1959-11-04','Paraskevi','Luby','F','1994-01-26'), (10103,'1953-11-26','Akemi','Birch','M','1986-12-02'), (10104,'1961-11-19','Xinyu','Warwick','M','1987-04-16'), (10105,'1962-02-05','Hironoby','Piveteau','M','1999-03-23'), (10106,'1952-08-29','Eben','Aingworth','M','1990-12-19'), (10107,'1956-06-13','Dung','Baca','F','1994-03-22'), (10108,'1952-04-07','Lunjin','Giveon','M','1986-10-02'), (10109,'1958-11-25','Mariusz','Prampolini','F','1993-06-16'), (10110,'1957-03-07','Xuejia','Ullian','F','1986-08-22'), (10111,'1963-08-29','Hugo','Rosis','F','1988-06-19'), (10112,'1963-08-13','Yuichiro','Swick','F','1985-10-08'), (10113,'1963-11-13','Jaewon','Syrzycki','M','1989-12-24'), (10114,'1957-02-16','Munir','Demeyer','F','1992-07-17'), (10115,'1964-12-25','Chikara','Rissland','M','1986-01-23'), (10116,'1955-08-26','Dayanand','Czap','F','1985-05-28'); 4. Doris 安装配置 这里我们以单机版为例 首先下载 Doris 1.1 release版本: https://doris.apache.org/downloads/downloads.html 解压到指定目录 tar zxvf apache-doris-1.1.0-bin.tar.gz -C doris-1.1 解压后的目录结构是这样: . ├── apache_hdfs_broker │ ├── bin │ ├── conf │ └── lib ├── be │ ├── bin │ ├── conf │ ├── lib │ ├── log │ ├── minidump │ ├── storage │ └── www ├── derby.log ├── fe │ ├── bin │ ├── conf │ ├── doris-meta │ ├── lib │ ├── log │ ├── plugins │ ├── spark-dpp │ ├── temp_dir │ └── webroot └── udf ├── include └── lib 配置fe和be cd doris-1.0 # 配置 fe.conf 和 be.conf,这两个文件分别在fe和be的conf目录下 打开这个 priority_networks 修改成自己的IP地址,注意这里是CIDR方式配置IP地址 例如我本地的IP是172.19.0.12,我的配置如下: priority_networks = 172.19.0.0/24 ​ ###### 在be.conf配置文件最后加上下面这个配置 disable_stream_load_2pc=false 注意这里默认只需要修改 fe.conf 和 be.conf 同样的上面这个配置就行了 默认fe元数据的目录在 fe/doris-meta 目录下 be的数据存储在 be/storage 目录下 启动 FE sh fe/bin/start_fe.sh --daemon 启动BE sh be/bin/start_be.sh --daemon MySQL命令行连接FE,这里新安装的Doris集群默认用户是root和admin,密码是空 mysql -uroot -P9030 -h127.0.0.1 Welcome to the MySQL monitor. Commands end with ; or \g. Your MySQL connection id is 41 Server version: 5.7.37 Doris version trunk-440ad03 ​ Copyright (c) 2000, 2022, Oracle and/or its affiliates. ​ Oracle is a registered trademark of Oracle Corporation and/or its affiliates. Other names may be trademarks of their respective owners. ​ Type 'help;' or '\h' for help. Type '\c' to clear the current input statement. ​ mysql> show frontends; +--------------------------------+-------------+-------------+----------+-----------+---------+----------+----------+------------+------+-------+-------------------+---------------------+----------+--------+---------------+------------------+ | Name | IP | EditLogPort | HttpPort | QueryPort | RpcPort | Role | IsMaster | ClusterId | Join | Alive | ReplayedJournalId | LastHeartbeat | IsHelper | ErrMsg | Version | CurrentConnected | +--------------------------------+-------------+-------------+----------+-----------+---------+----------+----------+------------+------+-------+-------------------+---------------------+----------+--------+---------------+------------------+ | 172.19.0.12_9010_1654681464955 | 172.19.0.12 | 9010 | 8030 | 9030 | 9020 | FOLLOWER | true | 1690644599 | true | true | 381106 | 2022-06-22 18:13:34 | true | | trunk-440ad03 | Yes | +--------------------------------+-------------+-------------+----------+-----------+---------+----------+----------+------------+------+-------+-------------------+---------------------+----------+--------+---------------+------------------+ 1 row in set (0.01 sec) ​ 将BE节点加入到集群中 mysql>alter system add backend "172.19.0.12:9050"; 这里是你自己的IP地址 查看BE mysql> show backends; +-----------+-----------------+-------------+---------------+--------+----------+----------+---------------------+---------------------+-------+----------------------+-----------------------+-----------+------------------+---------------+---------------+---------+----------------+--------------------------+--------+---------------+-------------------------------------------------------------------------------------------------------------------------------+ | BackendId | Cluster | IP | HeartbeatPort | BePort | HttpPort | BrpcPort | LastStartTime | LastHeartbeat | Alive | SystemDecommissioned | ClusterDecommissioned | TabletNum | DataUsedCapacity | AvailCapacity | TotalCapacity | UsedPct | MaxDiskUsedPct | Tag | ErrMsg | Version | Status | +-----------+-----------------+-------------+---------------+--------+----------+----------+---------------------+---------------------+-------+----------------------+-----------------------+-----------+------------------+---------------+---------------+---------+----------------+--------------------------+--------+---------------+-------------------------------------------------------------------------------------------------------------------------------+ | 10002 | default_cluster | 172.19.0.12 | 9050 | 9060 | 8040 | 8060 | 2022-06-22 12:51:58 | 2022-06-22 18:15:34 | true | false | false | 4369 | 328.686 MB | 144.083 GB | 196.735 GB | 26.76 % | 26.76 % | {"location" : "default"} | | trunk-440ad03 | {"lastSuccessReportTabletsTime":"2022-06-22 18:15:05","lastStreamLoadTime":-1,"isQueryDisabled":false,"isLoadDisabled":false} | +-----------+-----------------+-------------+---------------+--------+----------+----------+---------------------+---------------------+-------+----------------------+-----------------------+-----------+------------------+---------------+---------------+---------+----------------+--------------------------+--------+---------------+-------------------------------------------------------------------------------------------------------------------------------+ 1 row in set (0.00 sec) Doris单机版安装完成 5. Flink安装配置 5.1 下载安装Flink1.14.4 wget https://dlcdn.apache.org/flink/flink-1.14.4/flink-1.14.5-bin-scala_2.12.tgz tar zxvf flink-1.14.4-bin-scala_2.12.tgz 然后需要将下面的依赖拷贝到Flink安装目录下的lib目录下,具体的依赖的lib文件如下: wget https://jiafeng-1308700295.cos.ap-hongkong.myqcloud.com/flink-doris-connector-1.14_2.12-1.0.0-SNAPSHOT.jar wget https://repo1.maven.org/maven2/com/ververica/flink-sql-connector-mysql-cdc/2.2.1/flink-sql-connector-mysql-cdc-2.2.1.jar 启动Flink bin/start-cluster.sh 启动后的界面如下: 6. 开始同步数据到Doris 6.1 创建Doris数据库及表 create database demo; use demo; CREATE TABLE all_employees_info ( emp_no int NOT NULL, birth_date date, first_name varchar(20), last_name varchar(20), gender char(2), hire_date date, database_name varchar(50), table_name varchar(200) ) UNIQUE KEY(`emp_no`, `birth_date`) DISTRIBUTED BY HASH(`birth_date`) BUCKETS 1 PROPERTIES ( "replication_allocation" = "tag.location.default: 1" ); 6.2 进入Flink SQL Client bin/sql-client.sh embedded 开启 checkpoint,每隔10秒做一次 checkpoint Checkpoint 默认是不开启的,我们需要开启 Checkpoint 提交事务。 Source在启动时会扫描全表,将表按照主键分成多个chunk。并使用增量快照算法逐个读取每个chunk的数据。作业会周期性执行Checkpoint,记录下已经完成的chunk。当发生Failover时,只需要继续读取未完成的chunk。当chunk全部读取完后,会从之前获取的Binlog位点读取增量的变更记录。Flink作业会继续周期性执行Checkpoint,记录下Binlog位点,当作业发生Failover,便会从之前记录的Binlog位点继续处理,从而实现Exactly Once语义 SET execution.checkpointing.interval = 10s; 注意: 这里是演示,生产环境建议checkpoint间隔60秒 6.3 创建MySQL CDC表 在Flink SQL Client 下执行下面的 SQL CREATE TABLE employees_source ( database_name STRING METADATA VIRTUAL, table_name STRING METADATA VIRTUAL, emp_no int NOT NULL, birth_date date, first_name STRING, last_name STRING, gender STRING, hire_date date, PRIMARY KEY (`emp_no`) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = 'localhost', 'port' = '3306', 'username' = 'root', 'password' = 'MyNewPass4!', 'database-name' = 'emp_[0-9]+', 'table-name' = 'employees_[0-9]+' ); 'database-name' = 'emp_[0-9]+': 这里是使用了正则表达式,同时连接多个库 'table-name' = 'employees_[0-9]+':这里是使用了正则表达式,同时连接多个表 查询CDC表,我们可以看到下面的数据,标识一切正常 select * from employees_source limit 10; 6.4 创建 Doris Sink 表 CREATE TABLE cdc_doris_sink ( emp_no int , birth_date STRING, first_name STRING, last_name STRING, gender STRING, hire_date STRING, database_name STRING, table_name STRING ) WITH ( 'connector' = 'doris', 'fenodes' = '172.19.0.12:8030', 'table.identifier' = 'demo.all_employees_info', 'username' = 'root', 'password' = '', 'sink.properties.two_phase_commit'='true', 'sink.label-prefix'='doris_demo_emp_001' ); 参数说明: connector : 指定连接器是doris fenodes:doris FE节点IP地址及http port table.identifier : Doris对应的数据库及表名 username:doris用户名 password:doris用户密码 sink.properties.two_phase_commit:指定使用两阶段提交,这样在stream load的时候,会在http header里加上 two_phase_commit:true ,不然会失败 sink.label-prefix : 这个是在两阶段提交的时候必须要加的一个参数,才能保证两端数据一致性,否则会失败 其他参数参考官方文档 https://doris.apache.org/zh-CN/docs/ecosystem/flink-doris-connector.html 这个时候查询Doris sink表是没有数据的 select * from cdc_doris_sink; 6.5 将数据插入到Doris表里 执行下面的SQL: insert into cdc_doris_sink (emp_no,birth_date,first_name,last_name,gender,hire_date,database_name,table_name) select emp_no,cast(birth_date as string) as birth_date ,first_name,last_name,gender,cast(hire_date as string) as hire_date ,database_name,table_name from employees_source; 然后我们可以看到Flink WEB UI上的任务运行信息 这里我们可以看看TaskManager的日志信息,会发现这里是使用两阶段提交的,而且数据是通过http chunked方式不断朝BE端进行传输的,知道Checkpoint,才会停止。Checkpoint完成后会继续下一个任务的提交。 2022-06-22 19:04:01,350 INFO io.debezium.relational.history.DatabaseHistoryMetrics [] - Started database history recovery 2022-06-22 19:04:01,350 INFO io.debezium.relational.history.DatabaseHistoryMetrics [] - Finished database history recovery of 0 change(s) in 0 ms 2022-06-22 19:04:01,351 INFO io.debezium.util.Threads [] - Requested thread factory for connector MySqlConnector, id = mysql_binlog_source named = binlog-client 2022-06-22 19:04:01,352 INFO io.debezium.connector.mysql.MySqlStreamingChangeEventSource [] - Skip 0 events on streaming start 2022-06-22 19:04:01,352 INFO io.debezium.connector.mysql.MySqlStreamingChangeEventSource [] - Skip 0 rows on streaming start 2022-06-22 19:04:01,352 INFO io.debezium.util.Threads [] - Creating thread debezium-mysqlconnector-mysql_binlog_source-binlog-client 2022-06-22 19:04:01,374 INFO io.debezium.util.Threads [] - Creating thread debezium-mysqlconnector-mysql_binlog_source-binlog-client 2022-06-22 19:04:01,381 INFO io.debezium.connector.mysql.MySqlStreamingChangeEventSource [] - Connected to MySQL binlog at localhost:3306, starting at MySqlOffsetContext [sourceInfoSchema=Schema{io.debezium.connector.mysql.Source:STRUCT}, sourceInfo=SourceInfo [currentGtid=null, currentBinlogFilename=mysql_bin.000005, currentBinlogPosition=211725, currentRowNumber=0, serverId=0, sourceTime=null, threadId=-1, currentQuery=null, tableIds=[], databaseName=null], partition={server=mysql_binlog_source}, snapshotCompleted=false, transactionContext=TransactionContext [currentTransactionId=null, perTableEventCount={}, totalEventCount=0], restartGtidSet=null, currentGtidSet=null, restartBinlogFilename=mysql_bin.000005, restartBinlogPosition=211725, restartRowsToSkip=0, restartEventsToSkip=0, currentEventLengthInBytes=0, inTransaction=false, transactionId=null] 2022-06-22 19:04:01,381 INFO io.debezium.util.Threads [] - Creating thread debezium-mysqlconnector-mysql_binlog_source-binlog-client 2022-06-22 19:04:01,381 INFO io.debezium.connector.mysql.MySqlStreamingChangeEventSource [] - Waiting for keepalive thread to start 2022-06-22 19:04:01,497 INFO io.debezium.connector.mysql.MySqlStreamingChangeEventSource [] - Keepalive thread is running 2022-06-22 19:04:08,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:04:08,321 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6963, "Label": "doris_demo_001_0_1", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 634, "NumberLoadedRows": 634, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 35721, "LoadTimeMs": 9046, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9041, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:04:08,321 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:04:08,321 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_2 2022-06-22 19:04:08,321 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:04:08,325 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 1 2022-06-22 19:04:08,329 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6963] commit successfully." } 2022-06-22 19:04:18,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:04:18,310 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6964, "Label": "doris_demo_001_0_2", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 0, "NumberLoadedRows": 0, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 0, "LoadTimeMs": 9988, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9983, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:04:18,310 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:04:18,310 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_3 2022-06-22 19:04:18,310 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:04:18,312 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 2 2022-06-22 19:04:18,317 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6964] commit successfully." } 2022-06-22 19:04:28,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:04:28,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6965, "Label": "doris_demo_001_0_3", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 0, "NumberLoadedRows": 0, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 0, "LoadTimeMs": 9998, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9993, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:04:28,308 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:04:28,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_4 2022-06-22 19:04:28,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:04:28,311 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 3 2022-06-22 19:04:28,316 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6965] commit successfully." } 2022-06-22 19:04:38,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:04:38,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6966, "Label": "doris_demo_001_0_4", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 0, "NumberLoadedRows": 0, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 0, "LoadTimeMs": 9999, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9994, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:04:38,308 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:04:38,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_5 2022-06-22 19:04:38,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:04:38,311 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 4 2022-06-22 19:04:38,317 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6966] commit successfully." } 2022-06-22 19:04:48,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:04:48,310 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6967, "Label": "doris_demo_001_0_5", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 0, "NumberLoadedRows": 0, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 0, "LoadTimeMs": 10000, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9996, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:04:48,310 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:04:48,310 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_6 2022-06-22 19:04:48,310 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:04:48,312 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 5 2022-06-22 19:04:48,317 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6967] commit successfully." } 2022-06-22 19:04:58,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:04:58,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6968, "Label": "doris_demo_001_0_6", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 0, "NumberLoadedRows": 0, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 0, "LoadTimeMs": 9998, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9993, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:04:58,308 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:04:58,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_7 2022-06-22 19:04:58,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:04:58,311 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 6 2022-06-22 19:04:58,316 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6968] commit successfully." } 2022-06-22 19:05:08,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:05:08,309 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6969, "Label": "doris_demo_001_0_7", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 0, "NumberLoadedRows": 0, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 0, "LoadTimeMs": 9999, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9995, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:05:08,309 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:05:08,309 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_8 2022-06-22 19:05:08,309 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:05:08,311 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 7 2022-06-22 19:05:08,316 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6969] commit successfully." } 2022-06-22 19:05:18,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:05:18,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6970, "Label": "doris_demo_001_0_8", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 0, "NumberLoadedRows": 0, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 0, "LoadTimeMs": 9999, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9993, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:05:18,308 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:05:18,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_9 2022-06-22 19:05:18,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:05:18,311 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 8 2022-06-22 19:05:18,317 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6970] commit successfully." } 2022-06-22 19:05:28,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:05:28,310 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6971, "Label": "doris_demo_001_0_9", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 0, "NumberLoadedRows": 0, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 0, "LoadTimeMs": 10000, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9996, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:05:28,310 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:05:28,310 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_10 2022-06-22 19:05:28,310 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:05:28,315 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 9 2022-06-22 19:05:28,320 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6971] commit successfully." } 2022-06-22 19:05:38,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:05:38,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6972, "Label": "doris_demo_001_0_10", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 0, "NumberLoadedRows": 0, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 0, "LoadTimeMs": 9998, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 9992, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:05:38,308 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:05:38,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_11 2022-06-22 19:05:38,308 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:05:38,311 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 10 2022-06-22 19:05:38,316 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6972] commit successfully." } 2022-06-22 19:05:48,303 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load stopped. 2022-06-22 19:05:48,315 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - load Result { "TxnId": 6973, "Label": "doris_demo_001_0_11", "TwoPhaseCommit": "true", "Status": "Success", "Message": "OK", "NumberTotalRows": 520, "NumberLoadedRows": 520, "NumberFilteredRows": 0, "NumberUnselectedRows": 0, "LoadBytes": 29293, "LoadTimeMs": 10005, "BeginTxnTimeMs": 0, "StreamLoadPutTimeMs": 0, "ReadDataTimeMs": 0, "WriteDataTimeMs": 10001, "CommitAndPublishTimeMs": 0 } ​ 2022-06-22 19:05:48,315 INFO org.apache.doris.flink.sink.writer.RecordBuffer [] - start buffer data, read queue size 0, write queue size 3 2022-06-22 19:05:48,315 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - stream load started for doris_demo_001_0_12 2022-06-22 19:05:48,315 INFO org.apache.doris.flink.sink.writer.DorisStreamLoad [] - start execute load 2022-06-22 19:05:48,322 INFO org.apache.flink.streaming.runtime.operators.sink.AbstractStreamingCommitterHandler [] - Committing the state for checkpoint 11 2022-06-22 19:05:48,327 INFO org.apache.doris.flink.sink.committer.DorisCommitter [] - load result { "status": "Success", "msg": "transaction [6973] commit successfully." } 6.6 查询Doris 数据 这里我是插入了636条数据, mysql> select count(1) from all_employees_info ; +----------+ | count(1) | +----------+ | 634 | +----------+ 1 row in set (0.01 sec) ​ mysql> select * from all_employees_info limit 20; +--------+------------+------------+-------------+--------+------------+---------------+-------------+ | emp_no | birth_date | first_name | last_name | gender | hire_date | database_name | table_name | +--------+------------+------------+-------------+--------+------------+---------------+-------------+ | 10001 | 1953-09-02 | Georgi | Facello | M | 1986-06-26 | emp_1 | employees_1 | | 10002 | 1964-06-02 | Bezalel | Simmel | F | 1985-11-21 | emp_1 | employees_1 | | 10003 | 1959-12-03 | Parto | Bamford | M | 1986-08-28 | emp_1 | employees_1 | | 10004 | 1954-05-01 | Chirstian | Koblick | M | 1986-12-01 | emp_1 | employees_1 | | 10005 | 1955-01-21 | Kyoichi | Maliniak | M | 1989-09-12 | emp_1 | employees_1 | | 10006 | 1953-04-20 | Anneke | Preusig | F | 1989-06-02 | emp_1 | employees_1 | | 10007 | 1957-05-23 | Tzvetan | Zielinski | F | 1989-02-10 | emp_1 | employees_1 | | 10008 | 1958-02-19 | Saniya | Kalloufi | M | 1994-09-15 | emp_1 | employees_1 | | 10009 | 1952-04-19 | Sumant | Peac | F | 1985-02-18 | emp_1 | employees_1 | | 10010 | 1963-06-01 | Duangkaew | Piveteau | F | 1989-08-24 | emp_1 | employees_1 | | 10011 | 1953-11-07 | Mary | Sluis | F | 1990-01-22 | emp_1 | employees_1 | | 10012 | 1960-10-04 | Patricio | Bridgland | M | 1992-12-18 | emp_1 | employees_1 | | 10013 | 1963-06-07 | Eberhardt | Terkki | M | 1985-10-20 | emp_1 | employees_1 | | 10014 | 1956-02-12 | Berni | Genin | M | 1987-03-11 | emp_1 | employees_1 | | 10015 | 1959-08-19 | Guoxiang | Nooteboom | M | 1987-07-02 | emp_1 | employees_1 | | 10016 | 1961-05-02 | Kazuhito | Cappelletti | M | 1995-01-27 | emp_1 | employees_1 | | 10017 | 1958-07-06 | Cristinel | Bouloucos | F | 1993-08-03 | emp_1 | employees_1 | | 10018 | 1954-06-19 | Kazuhide | Peha | F | 1987-04-03 | emp_1 | employees_1 | | 10019 | 1953-01-23 | Lillian | Haddadi | M | 1999-04-30 | emp_1 | employees_1 | | 10020 | 1952-12-24 | Mayuko | Warwick | M | 1991-01-26 | emp_1 | employees_1 | +--------+------------+------------+-------------+--------+------------+---------------+-------------+ 20 rows in set (0.00 sec) 6.7 测试删除 mysql> use emp_2; Reading table information for completion of table and column names You can turn off this feature to get a quicker startup with -A ​ Database changed mysql> show tables; +-----------------+ | Tables_in_emp_2 | +-----------------+ | employees_1 | | employees_2 | +-----------------+ 2 rows in set (0.00 sec) ​ mysql> delete from employees_2 where emp_no in (12013,12014,12015); Query OK, 3 rows affected (0.01 sec) 验证Doris数据删除 mysql> select count(1) from all_employees_info ; +----------+ | count(1) | +----------+ | 631 | +----------+ 1 row in set (0.01 sec) 7. 总结 本问主要介绍了FLink CDC分库分表怎么实时同步,并结合Apache Doris Flink Connector最新版本整合的Flink 2PC 和 Doris Stream Load 2PC的机制及整合原理,使用方法等。 希望能给大家带来一点帮助。 8.相关链接: SelectDB 官方网站: https://selectdb.com Apache Doris 官方网站: http://doris.apache.org Apache Doris Github: https://github.com/apache/doris Apache Doris 开发者邮件组: dev@doris.apache.org

优秀的个人博客,低调大师

工信部刘烈宏:提速降费成果显著,继续实施精准降费

4 月 19 日消息 根据中国工信部官方公众号 @工信微报消息,今日国务院新闻办举行网络提速降费政策有关情况政策例行吹风会。工业和信息化部党组成员、副部长刘烈宏,工业和信息化部信息通信发展司副司长刘郁林,国务院国资委财管运行局负责人刘绍娓出席吹风会介绍有关情况,并回答记者提问。 刘烈宏表示,我国提速降费政策实施以来,取得了显著成就,实现了农村和城市“同网同速”。自 2015 年网络提速降费实施以来,这五年间固定宽带单位带宽和移动网络单位流量平均资费降幅超过 95%,2020 年下半年以来,随着 5G 建设发展进程加速,移动网络单位流量平均资费又下降了 10% 以上。我国移动通信用户月均流量(DOU)从 2015 年初的 205MB 提升至 10.85GB,提升 40 多倍。 刘烈宏还表示,将“光纤入户”纳入城镇老旧小区改造内容,保障网络建设通行权。开展学校联网攻坚,全国中小学校(含教学点)100% 实现宽带接入。推动远程医疗能力覆盖所有贫困县县级医院。 今年我国《政府工作报告》提出,加大 5G 网络和千兆光网建设力度,丰富应用场景,中小企业宽带和专线平均资费再降 10%。具体来看,目前我国基础电信企业面向中小企业开展企业宽带提速惠企行动,推出“云 + 网 + 应用”等融合产品优惠,持续降低了企业用网用云成本,助力中小企业提升信息化水平。还面向农村的脱贫户提出 5 折基础通信服务资费折扣,面向听障人群推出畅听王卡;面向老年人推出银龄卡、孝心卡等。 农村方面,工信部表示,全国行政村通光纤、通 4G 比例都超过 99%,经测试,已通光纤试点村平均下载速率超过 100Mb/s。今年再部署第 7 批电信普遍服务的建设任务,预计在农村及偏远地区支持 1 万个 4G 基站建设,推动宽带网络逐步向农村人口聚居区、生产作业区、交通要道沿线等重点区域延伸。

优秀的个人博客,低调大师

【云栖号案例 | 互联网】华大基因:打造精准医疗应用云平台日志方案

云栖号案例库:【点击查看更多上云案例】不知道怎么上云?看云栖号案例库,了解不同行业不同发展阶段的上云方案,助力你上云决策! 公司简介 华大基因是中国最领先的基因科技公司,华大基因为消除人类病痛、经济危机、国家灾难、濒危动物保护、缩小贫富差距等方面提供分子遗传层面的技术支持。目前,世界上只有两个国家的三个公司可以生产、量产临床级别的基因测序仪,华大基因是中国的唯一一家。我们在基因的产权研发方面从1999年开始做了很多的工作。在2014年,我们与阿里云有了初步的接触,在2015年上线了我国第一个基因云计算平台。 业务痛点 我们与阿里云合作是因为我们看到基因技术从过去的只在实验室中逐渐进入到广大群众的生活场景当中,不管是在医学健康方面、生殖健康方面、肿瘤防治方面、病原感染方面还是农业育种,以及与我们每个人息息相关的健康管理,基因技术已经取得越来越多的应用场景,在国产基因测序仪的助力之下,基因数据产生的体量也越来越庞大,远远的超出了原有的计算能力所能支持的范围。 解决方案 针对上述情况,华大基因业务逐步迁移到阿里云计算平台之上。 新的日志分析架构如页面下方架构图所示。 计算集群:本地IDC作为原始测序数据(FastQ)的计算集群。 存储:阿里云OSS用于比对结果数据和测序数据。在存储方面,我们也使用了阿里的产品,每年我们会产生非常多的基因数据,明年我们计划对十万人进行基因组的基因测序和分析,我们将与阿里云计算平台一起在2018年用国产测序仪完成计算、分析和交付。 大数据计算:批量计算、Maxcompute等一些异构计算方式,使我们原先需要几周甚至更长时间才能完成的计算任务在一两天内得以解决。在我们现在进行的百万人基因组项目中,阿里云的Maxcompute技术帮助我们大大加速了对于人群结构的分析速度的进展。 1.使用阿里云MaxCompute处理群体变异检测和人群遗传结构分析。2.通过BatchCompute完成数据质控和比对。 上云价值 另外,在对百万人的基因数据进行遗传结构分析时,我们需要把每一个人与剩余的所有人进行遗传距离计算,这个计算量是巨大的,计算复杂度已经远远超出了传统计算条件下硬件设备所能承受的能力范围,通过使用Maxcompute,我们已经在这方面取得了技术突破,其中,我们在几小时内就可以把一个人与十万人中所有遗传距离进行计算,计算成本大幅降低至1000美金以内,这样的例子我们还在不断的开发中,相信Maxcompute也会给我们带来更多的惊喜。 相关产品 大数据计算服务 · MaxCompute MaxCompute(原ODPS)是一项大数据计算服务,它能提供快速、完全托管的PB级数据仓库解决方案,使您可以经济并高效的分析处理海量数据。更多关于阿里云MaxCompute的介绍,参见MaxCompute产品详情页。 批量计算 批量计算(BatchCompute)是一种适用于大规模并行批处理作业的分布式云服务。BatchCompute可支持海量作业并发规模,系统自动完成资源管理,作业调度和数据加载,并按实际使用量计费。BatchCompute广泛应用于电影动画渲染、生物数据分析、多媒体转码、金融保险分析、科学计算等领域。 更多关于批量计算的介绍,参见批量计算产品详情页。 对象存储OSS 阿里云对象存储服务(Object Storage Service,简称 OSS),是阿里云提供的海量、安全、低成本、高可靠的云存储服务。其数据设计持久性不低于 99.9999999999%(12 个 9),服务设计可用性(或业务连续性)不低于 99.995%。 更多关于对象存储OSS的介绍,参见对象存储OSS产品详情页。 【云栖号在线课堂】每天都有产品技术专家分享! 在线课堂地址:https://yqh.aliyun.com/zhibo 立即加入社群,与专家面对面,及时了解课程最新动态! 【云栖号在线课堂 社群】https://c.tb.cn/F3.Z8gvnK

优秀的个人博客,低调大师

看IBM如何为汽车信息化精准定位

2017年中国汽车CIO峰会在对百余位汽车企业CIO展开书面问卷调查后,分析出近两年为数众多的汽车厂商在进行企业数字化转型过程中,都在部署物联网、云计算、大数据分析、移动办公等解决方案,同时指出转型聚焦化困难,缺少完整数字化蓝图,对如何更好的进行数据分析、物联网服务、云端建设、企业IT架构建模等方面寻求数字化转型的顶层设计方案。 对于数字信息化建设,ACS2017组委会得到全球最大的信息技术和业务解决方案公司IBM的全力支持,将于10月26日在大数据、云计算、车联网、移动信息化专场为我们带来精彩报告分享! 在信息化时代背景下,IBM于去年宣布“认知商业”战略,并确信这一转型将为更多的企业塑造数字化的未来。在基于云计算、大数据分析和物联网等新兴技术商业模式下,IBM的认知技术将以全新方式将数据连结在一起,获得全新洞察。全方位布局物联网服务,以“泛连接、云平台、轻应用、大数据”相结合为趋势,在汽车与互联网的融合上实现技术突破。 对汽车企业CIO来说,实现信息化建设进程中,物联网公有云建设技术分享是共同期待的话题。通过物联网技术解决了生产制造过程中哪些业务痛点,例如定位,设备诊断,传感器等方向;通过物联网技术增加制造企业从控制层,生产制造、产品、消费者之间的距离,使制造服务直接面向用户等。 开发认知型车联网,让认知计算技术与物联网云平台相结合,在推动汽车行业转型的同时亦能加快人工智能的商业化进程。IBM已与多家知名汽车企业就车载技术、物联网数据方面进行合作,研发联网汽车服务产品,开发新型数字化汽车服务,满足车联网环境下的汽车消费者的新兴需求。在智能工厂、智能供应链、数字化流程、虚拟工厂等方向正在展开。 汽车企业面对的市场竞争非常激烈,市场环境瞬息万变,IBM 助力汽车行业借助信息技术专业洞察,结合全面的技术和经济有效的方法,可帮助车企实现精益生产、转变其价值链,为汽车行业提供热门解决方案。 基于四大模块: -智慧研发体系:在产品生命周期管理、系统工程、产品组合管理与项目管理、工程设计云上推动业务转型。 -市场营销与服务:结合数据分析,推动汽车业电子商务、业务分析及优化能力。 -精益生产与供应链管理:超越传统ERP,推动业务转型、面向汽车行业整合服务解决方案,推动复杂生产系统向智能化的设备制造转变、建立长效的“零”停台管理机制,提升 IT 支持能力。 -创新增值业务:基于车联网、云计算、移动社交方式推进全新的业务模式,颠覆传统行业结构,更直接地融合到汽车企业战略和举措中。 IBM始终关注各行业企业的信息化建设及数字化转型,并致力于为提供全球最领先的人工智能技术,此次,中国汽车CIO峰会以“互联网+时代的汽车全产业链信息化解决方案”为主题,携手350+汽车企业高层、信息技术专家共同为您呈现一场信息化行业盛宴。在论坛上我们将着力为您解读信息技术与智能化应用在汽车行业的发展特点,突出“互联网+”为标志的数字化转型趋势,努力在汽车信息化领域带来更多技术革新,完成数字化变革。 距离大会召开倒计时只剩一个月,目前,大会已邀请到100+汽车整车、零部件、经销商CEO、CIO、CTO、IT负责人参与。包括:北汽集团信息技术部部长李晓龙、宝马集团CIO陈宇星、上汽集团信息战略和系统支持部金忠孝、福特汽车亚太智能移动IT总监侯新海、北汽福田副总杨国涛、江淮汽车副总工程师李世杭、吉利汽车(战略规划部)总监刘景林、苏州金龙IT部长吴震、前途汽车信息总监邢红波、比亚迪CIO裘彦近40+整车企业。零部件企业:马瑞利亚太区CIO方志强、赛轮金宇CIO朱小兵、福耀玻璃IT副总裁夏乐冰、万丰奥特CIO汪清跃、万向钱潮IT经理刘华、中鼎控股智能化研究中心主任李科伍、一汽富维信息管理室主任陈霖等近70+零部件企业信息化负责人共同交流,大会现场更有“一对一”VIP对接,为您打造更为高效的沟通平台!如果您想接触更多具体参考价值的案例分享,机会不容错过!期待您的参会! 本文出处:畅享网 本文来自云栖社区合作伙伴畅享网,了解相关信息可以关注vsharing.com网站。

优秀的个人博客,低调大师

风控开发指南:Go集成查行政处罚实现精准合规审查

在大型B2B采购平台与分布式供应链管理系统中,确保入驻企业的资质健康与经营合规是构建高可用业务基座的必要条件。传统的企业审核往往依赖商家手动上传营业执照、人工比对各类行政公开信息。这种方式不仅流程繁琐、并发处理能力极弱,而且存在严重的数据滞后与信息不匹配隐患,无法满足海量商家高并发入驻的实时合规审查需求。

资源下载

更多资源
Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册