MatrixOrigin MatrixOrigin 文档
产品文档
MatrixOne 当前产品 MatrixOne Intelligence
关于 MatrixOne 快速开始 开发指南 教程 部署指南 运维 数据迁移 测试 性能调优 安全与权限 参考手册 故障诊断 常见问题解答 版本发布纪要 名词术语表 社区贡献指南
/
Contents Menu Expand Light mode Dark mode Auto light/dark, in light mode Auto light/dark, in dark mode Skip to content
MatrixOne 文档
MatrixOne 文档
  • 主页
  • 关于 MatrixOne
    • MatrixOne 功能清单
    • MatrixOne 功能概述
      • Git for Data
      • 多租户
      • 极致扩展性
      • 高性价比
      • 高可用
      • 时序
      • 流
      • 用户定义函数
      • MySQL 兼容性
      • feature-overview
    • MatrixOne 技术架构
      • 存储引擎架构详解
      • Logservice 架构详解
      • Logtail 协议详解
      • 事务与锁机制实现详解
      • Proxy 架构详解
      • WAL 技术详解
      • 数据缓存及冷热数据分离架构详解
      • 流引擎架构详解
      • MatrixOne-Operator 设计与实现详解
    • MatrixOne 与其它数据库对比
      • MatrixOne与常见OLTP数据库的对比
    • 最新动态
  • 快速开始
    • 在 macOS 上部署
      • 使用二进制包部署
      • 使用 Docker 部署
    • 在 Linux 上部署
      • 使用二进制包部署
      • 使用 Docker 部署
    • SQL 的基本操作
  • 开发指南
    • Java 连接
      • Java ORMs 连接
    • C# 连接
    • Python 连接
    • Golang 连接
    • TypeScript 连接
    • configure-mo-ssl-connection
    • ODBC 连接
    • Power BI 连接
    • 数据库模式设计
      • 创建数据库
      • 创建表
      • 复制表
      • 创建视图
      • 创建临时表
      • 创建次级索引
      • 创建约束
        • NOT NULL 非空约束
        • UNIQUE KEY 唯一约束
        • PRIMARY KEY 主键约束
        • FOREIGN KEY 外键约束
        • AUTO INCREMENT 自增约束
      • 1.1-overview
    • 数据写入
      • 批量插入
        • 插入 csv 文件
        • 插入 jsonlines 文件
        • 从对象存储导入文件
        • Source 插入
      • 更新数据
      • 删除数据
      • 预处理
      • stream-load
    • 数据写出
      • mo-dump 工具写出
    • 数据读取
      • 多表连接查询
      • 子查询
      • 视图
      • 公共表表达式
      • 窗口函数
        • 时间窗口
    • 数据去重
      • BITMAP
    • 数据集成
    • 租户设计
      • 发布订阅
    • 事务
      • MatrixOne 的事务
        • 显式事务
        • 隐式事务
        • 悲观事务
        • 乐观事务
        • 隔离级别
        • MVCC
        • 使用指南
          • 应用场景
        • 应用场景
    • 用户定义函数
      • UDF python 进阶
    • 向量
      • 向量检索
      • IVF 向量查询排名选项
      • 聚类中心
      • vector_index
    • AI Agent 工具
      • 使用 AI Agent 查询 MatrixOne 文档
    • 应用开发示例
      • Python 基础示例
      • SpringBoot 和 JPA 基础示例
      • SpringBoot 和 MyBatis 基础示例
      • SQLAlchemy 基础示例
      • Django 基础示例
      • Golang 基础示例
      • Gorm 基础示例
      • C# 基础示例
      • Typescript 基础示例
      • HTAP 应用基础示例
      • RAG 应用基础示例
      • 以图(文)搜图应用基础示例
      • Dify 平台接入 MatrixOne 指南
      • 开始前准备
      • 操作步骤
      • git4data-demo
      • MatrixOne Python SDK 适配示例
      • 开始前准备
      • 快速入门
    • 生态工具
      • BI 工具
        • 通过永洪 BI 实现 MatrixOne 的可视化报表
        • 通过 Superset 实现 MatrixOne 可视化监控
      • ETL 工具
        • 从 MySQL 写入数据到 MatrixOne
        • 从 Oracle 写入数据到 MatrixOne
        • 使用 DataX 将数据写入 MatrixOne
          • 从 MySQL 写入数据到 MatrixOne
          • 从 Oracle 写入数据到 MatrixOne
          • 从 PostgreSQL 写入数据到 MatrixOne
          • 从 SQL Server 写入数据到 MatrixOne
          • 从 MongoDB 写入数据到 MatrixOne
          • 从 TiDB 写入数据到 MatrixOne
          • 从 ClickHouse 写入数据到 MatrixOne
          • 从 Doris 写入数据到 MatrixOne
          • 从 InfluxDB 写入数据到 MatrixOne
          • 从 Elasticsearch 写入数据到 MatrixOne
      • 计算引擎
        • 从 MySQL 写入数据到 MatrixOne
        • 从 Hive 写入数据到 MatrixOne
        • 从 Doris 写入数据到 MatrixOne
        • 使用 Flink 将实时数据写入 MatrixOne
          • 从 MySQL 写入数据到MatrixOne
          • 从 Oracle 写入数据到MatrixOne
          • 从 SQL Server 写入数据到MatrixOne
          • 从 PostgreSQL 写入数据到MatrixOne
          • 从 MongoDB 写入数据到MatrixOne
          • 从 TiDB 写入数据到MatrixOne
          • 从 Kafka 写入数据到MatrixOne
      • 调度工具
    • develop-overview
  • 教程
    • Python 基础示例
    • SpringBoot 和 JPA 基础示例
    • SpringBoot 和 MyBatis 基础示例
    • SQLAlchemy 基础示例
    • Django 基础示例
    • Golang 基础示例
    • Gorm 基础示例
    • C# 基础示例
    • Typescript 基础示例
    • HTAP 应用基础示例
    • RAG 应用基础示例
    • 以图(文)搜图应用基础示例
    • Dify 平台接入 MatrixOne 指南
    • 开始前准备
    • 操作步骤
    • git4data-demo
    • MatrixOne Python SDK 适配示例
    • 开始前准备
    • 快速入门
  • 部署指南
    • 集群拓扑规划
      • 体验环境
      • 最小生产环境
      • 推荐生产环境
    • 集群部署指南
      • 已部署 Kubernetes 和对象存储环境
    • 集群运维管理
      • 版本升级
      • 健康检查与资源监控
      • 集群扩缩容
      • 负载与租户隔离
      • 本地对象存储导入数据
      • Operator 管理
  • 运维
    • 备份与恢复相关概念
    • mo-dump 备份与恢复
    • mo_br 备份与恢复
      • 原理概述
      • 示例
      • mo_br snapshot
      • mo_br pitr
    • MatrixOne 主备容灾
    • 数据变更捕获
      • MatrixOne 到 MySQL
      • MatrixOne 到 MatrixOne
    • 挂载目录到 Docker 容器
  • 数据迁移
    • 将数据从 MySQL 迁移至 MatrixOne
    • 将数据从 Oracle 迁移至 MatrixOne
    • 将数据从 SQL Server 迁移至 MatrixOne
    • 将数据从 PostgreSQL 迁移至 MatrixOne
  • 测试
    • TPCH 测试
    • TPCC 测试
    • 测试工具
      • MO-Tester 规范要求
  • 性能调优
    • MatrixOne 执行计划
      • 使用 EXPLAIN 理解执行计划
      • JOIN 查询的执行计划
      • 子查询的执行计划
      • 聚合查询的执行计划
      • 视图的执行计划
    • 性能调优最佳实践
      • 通过扩展 CN 提升性能
      • through-index
      • 使用分区表调优
        • 分区裁剪
        • 分区裁剪在 KEY 分区表上的应用
        • 分区裁剪在 Hash 分区表上的应用
        • 分区剪裁的性能调优示例
        • 限制
    • optimizer-hints
  • 安全与权限
    • 身份鉴别与认证
    • 密码管理
    • 权限管理
      • 场景案例
      • 最佳实践
      • 操作指南
        • 创建租户,验证资源隔离
        • 新租户创建用户、创建角色和授权
    • 操作指南
      • 创建租户,验证资源隔离
      • 新租户创建用户、创建角色和授权
    • 数据加密传输
    • 安全审计
  • 参考手册
    • 系统变量参数
      • 保存查询结果支持
      • 时区支持
      • 大小写敏感支持
      • 外键检查支持
      • 查询结果集列名与用户指定大小写一致支持
      • 非法登录限制
      • 密码复杂度校验
      • 连接白名单
      • enable-remap-hint
    • cte_max_memory_bytes
    • cte_max_recursion_depth
    • enable_explain_scheduling
    • event_scheduler
    • experimental_cagra_index
    • experimental_fulltext_index
    • experimental_hnsw_index
    • experimental_ivf_index
    • experimental_ivfpq_index
    • fulltext_bloom_filter_pushdown
    • group_concat_max_len
    • lock_wait_timeout
    • protected_databases
    • query_max_workers
    • query_pool_strict
    • remap_rewrites
    • remapdb
    • sort_spill_mem
    • 自定义变量
    • SQL 结构与语法
      • 注释
    • 数据类型
      • 数据类型转换
      • 日期和时间类型
        • YEAR 类型
      • Geometry 数据类型
      • 精确数值类型-DECIMAL
      • 向量数据类型
      • BLOB 和 TEXT 数据类型
      • DATALINK 数据类型
      • ENUM 类型
      • JSON 数据类型
      • UUID 数据类型
      • set-type
    • SQL 目录
      • 数据定义语言(DDL)
        • CREATE INDEX
        • CREATE INDEX...USING IVFFLAT
        • CREATE INDEX...USING HNSW
        • CREATE FULLTEXT INDEX
        • CREATE TABLE
        • CREATE TABLE AS SELECT
        • CREATE TABLE ... LIKE
        • CREATE EXTERNAL TABLE
        • CREATE EXTERNAL TABLE ... ENGINE = DATASTREAM
        • CREATE ICEBERG CATALOG
        • ALTER ICEBERG CATALOG
        • DROP ICEBERG CATALOG
        • CREATE MONGODB CONNECTION
        • ALTER MONGODB CONNECTION
        • DROP MONGODB CONNECTION
        • CREATE CLONE
        • CREATE CLUSTER TABLE
        • CREATE PITR
        • CREATE PUBLICATION
        • CREATE SEQUENCE
        • CREATE STAGE
        • CREATE...FROM...PUBLICATION...
        • CREATE VIEW
        • CREATE FUNCTION...LANGUAGE SQL AS
        • CREATE FUNCTION...LANGUAGE PYTHON AS
        • CREATE SOURCE
        • CREATE DYNAMIC TABLE
        • CREATE SNAPSHOT
        • CREATE BRANCH
        • DELETE BRANCH
        • DIFF BRANCH
        • MERGE BRANCH
        • ALTER TABLE
        • ALTER TABLE ... ALTER REINDEX
        • ALTER PITR
        • ALTER PUBLICATION
        • ALTER SEQUENCE
        • ALTER STAGE
        • ALTER VIEW
        • DROP DATABASE
        • DROP INDEX
        • DROP TABLE
        • DROP PITR
        • DROP PUBLICATION
        • DROP SEQUENCE
        • DROP STAGE
        • DROP SNAPSHOT
        • DROP VIEW
        • DROP FUNCTION
        • TRUNCATE TABLE
        • RENAME TABLE
        • RESTORE PITR
        • RESTORE SNAPSHOT
        • branch-protect-snapshots
        • create-replace-view
        • data-branch-pick
        • sql-task
        • 数据分支权限
      • 数据修改语言(DML)
        • INSERT INTO SELECT
        • DELETE
        • UPDATE
        • LOAD DATA INFILE
        • LOAD DATA INLINE
        • UPSERT
          • INSERT ON DUPLICATE KEY UPDATE
          • INSERT IGNORE
          • REPLACE
        • MERGE
        • RETURNING
        • Information Functions
          • LAST_INSERT_ID()
          • ROW_COUNT()
        • case
        • Replace
      • 数据查询语言(DQL)
        • OUTER APPLY
        • JOIN
          • INNER JOIN
          • LEFT JOIN
          • RIGHT JOIN
          • FULL JOIN
          • OUTER JOIN
          • 示例
          • NATURAL JOIN
          • cross-join
        • SELECT
        • SUBQUERY
          • Derived Tables
          • 子查询与比较操作符的使用
          • SUBQUERY with ANY or SOME
          • SUBQUERY with ALL
          • SUBQUERY with EXISTS
          • SUBQUERY with IN
        • With CTE
        • BY RANK WITH OPTION
        • 联合查询概述
          • UNION
          • INTERSECT
          • MINUS
        • UNION
        • INTERSECT
        • MINUS
      • 数据控制语言(DCL)
        • ALTER ACCOUNT
        • CREATE ROLE
        • CREATE USER
        • ALTER USER
        • DROP ACCOUNT
        • DROP USER
        • DROP ROLE
        • GRANT
        • REVOKE
        • role-rule
      • 其他
        • SHOW DATABASES
        • SHOW CREATE TABLE
        • SHOW CREATE VIEW
        • SHOW CREATE PUBLICATION
        • SHOW TABLES
        • SHOW INDEX
        • SHOW COLLATION
        • SHOW COLUMNS
        • SHOW FUNCTION STATUS
        • SHOW GRANT
        • SHOW PROCESSLIST
        • SHOW PUBLICATIONS
        • SHOW ROLES
        • SHOW SEQUENCES
        • SHOW STAGES
        • SHOW SUBSCRIPTIONS
        • SHOW VARIABLES
        • SHOW ICEBERG CATALOGS
        • SHOW ICEBERG NAMESPACES
        • SHOW ICEBERG TABLES
        • SHOW MONGODB CONNECTIONS
        • SHOW PROFILE
        • show-pitr
        • show-snapshots
        • SET
        • USE DATABASE
          • ANALYZE TABLE
          • DUMP TABLE
          • LOAD TABLE
          • PERFORM
        • KILL
        • Prepared
          • EXECUTE
          • DEALLOCATE
        • Explain
          • EXPLAIN Output Format
          • Explain Analyze
          • Explain Prepared
        • Partition
        • ANALYZE TABLE
        • DUMP TABLE
        • LOAD TABLE
        • PERFORM
    • 运算符
      • Operators
        • 运算符的优先级
        • 算数运算符
          • %,MOD
          • *
          • +
          • -
          • -
          • /
          • DIV
        • 赋值运算符
          • =
        • 二进制运算符
          • &
          • >>
          • <<
          • ^
          • |
          • ~
        • 强制转换函数和运算符
          • BINARY
          • CAST
          • CONVERT
          • DECODE
          • ENCODE
          • SERIAL
          • SERIAL_FULL
        • 比较函数和运算符
          • >
          • >=
          • <
          • <>,!=
          • <=
          • =
          • <=>
          • BETWEEN ... AND ...
          • IN
          • IS
          • IS NOT
          • IS NOT NULL
          • IS NULL
          • ISNULL
          • ILIKE
          • LIKE
          • NOT BETWEEN ... AND ...
          • NOT IN
          • NOT LIKE
          • COALESCE
          • function_interval
          • function_greatest
          • function_least
          • function_strcmp
        • 控制流函数
          • CASE WHEN
          • IF
          • IFNULL
          • NULLIF
        • 逻辑运算符
          • AND,&&
          • NOT,!
          • OR
          • XOR
      • matrixone-function-list
        • 聚合函数
          • AVG
          • APPROX_PERCENTILE
          • BITMAP
          • BIT_AND
          • BIT_OR
          • BIT_XOR
          • COUNT
          • GROUP_CONCAT
          • HLL_ADD_AGG
          • HLL_CARDINALITY
          • HLL_MERGE_AGG
          • MAX
          • MAX_BY
          • MAX_BY_NON_NULL
          • MEDIAN
          • MIN
          • STDDEV_POP
          • SUM
          • VARIANCE
          • VAR_POP
        • 日期时间
          • CURDATE()
          • CURRENT_TIMESTAMP()
          • DATE()
          • DATE_ADD()
          • DATE_FORMAT()
          • DATE_SUB()
          • DATE_TRUNC()
          • DATEDIFF()
          • DAY()
          • DAYOFYEAR()
          • EXTRACT()
          • HOUR()
          • FROM_UNIXTIME
          • MINUTE()
          • MONTH()
          • NOW()
          • SECOND()
          • STR_TO_DATE()
          • SYSDATE()
          • TIME()
          • TIMEDIFF()
          • TIMESTAMP()
          • TIMESTAMPDIFF()
          • TO_DATE()
          • TO_DAYS()
          • TO_SECONDS()
          • UNIX_TIMESTAMP
          • UTC_TIMESTAMP()
          • WEEK()
          • WEEKDAY()
          • YEAR()
          • addtime
          • curtime
          • dayname
          • get-format
          • maketime
          • monthname
          • quarter
          • subtime
          • time-format
          • timestampadd
          • yearweek
        • Geo 函数
          • H3 Index Functions
          • MBRContains
          • MBRCoveredBy
          • MBRCovers
          • MBRDisjoint
          • MBREquals
          • MBRIntersects
          • MBROverlaps
          • MBRTouches
          • MBRWithin
          • S2 Cell Functions
          • ST_Area()
          • ST_AsGeoJSON
          • ST_AsText()
          • ST_AsWKB()
          • ST_Boundary()
          • ST_Buffer()
          • ST_Centroid()
          • ST_Collect()
          • ST_Contains()
          • ST_ConvexHull()
          • ST_CoveredBy()
          • ST_Covers()
          • ST_Crosses()
          • ST_Difference()
          • ST_Dimension()
          • ST_Disjoint()
          • ST_Distance()
          • ST_Distance_Sphere()
          • ST_EndPoint()
          • ST_Envelope()
          • ST_Equals()
          • ST_ExteriorRing()
          • ST_FrechetDistance()
          • ST_GeoHash
          • ST_GeomCollFromText()
          • ST_GeomCollFromWKB()
          • ST_GeometryN()
          • ST_GeometryType()
          • ST_GeomFromGeoJSON
          • ST_GeomFromText()
          • ST_GeomFromWKB()
          • ST_HausdorffDistance()
          • ST_InteriorRingN()
          • ST_Intersection()
          • ST_Intersects()
          • ST_IsClosed()
          • ST_IsCollection()
          • ST_IsEmpty()
          • ST_IsRing()
          • ST_IsSimple()
          • ST_IsValid()
          • ST_LatFromGeoHash
          • ST_Latitude()
          • ST_Length()
          • ST_LineFromText()
          • ST_LineFromWKB()
          • ST_LineInterpolatePoint()
          • ST_LineInterpolatePoints()
          • ST_LongFromGeoHash
          • ST_Longitude()
          • ST_MakeEnvelope()
          • ST_MLineFromText()
          • ST_MLineFromWKB()
          • ST_MPointFromText()
          • ST_MPointFromWKB()
          • ST_MPolyFromText()
          • ST_MPolyFromWKB()
          • ST_NumGeometries()
          • ST_NumInteriorRings()
          • ST_NumPoints()
          • ST_Overlaps()
          • ST_PointAtDistance()
          • ST_PointFromGeoHash
          • ST_PointFromText()
          • ST_PointFromWKB()
          • ST_PointN()
          • ST_PointOnSurface()
          • ST_PolyFromText()
          • ST_PolyFromWKB()
          • ST_Simplify()
          • ST_SRID()
          • ST_StartPoint()
          • ST_SwapXY()
          • ST_SymDifference()
          • ST_Touches()
          • ST_Union()
          • ST_Validate()
          • ST_Within()
          • ST_X()
          • ST_Y()
        • 数学函数
          • ACOS()
          • ATAN()
          • BIT_COUNT()
          • CEIL()
          • CEILING()
          • COS()
          • COT()
          • CRC32()
          • EXP()
          • FLOOR()
          • LN()
          • LOG()
          • LOG2()
          • LOG10()
          • PI()
          • POWER()
          • ROUND()
          • RAND()
          • SIN()
          • SINH()
          • TAN()
          • atan2
          • degrees
          • radians
          • sign
          • truncate
        • 字符串函数
          • BIT_LENGTH()
          • CHAR_LENGTH()
          • CONCAT()
          • CONCAT_WS()
          • EMPTY()
          • ENDSWITH()
          • FIELD()
          • FIND_IN_SET()
          • FORMAT()
          • FROM_BASE64()
          • HEX()
          • INSTR()
          • LCASE()
          • LEFT()
          • LENGTH()
          • LOCATE()
          • LOWER()
          • LPAD()
          • LTRIM()
          • MD5()
          • NAME_CONST()
          • OCT()
          • REPEAT()
          • REVERSE()
          • RPAD()
          • RTRIM()
          • SHA1()/SHA()
          • SHA2()
          • SPACE()
          • SPLIT_PART()
          • STARTSWITH()
          • STRCMP()
          • SUBSTRING()
          • SUBSTRING_INDEX()
          • TO_BASE64()
          • TRIM()
          • UCASE()
          • UNHEX()
          • UPPER()
          • 正则表达式
            • NOT REGEXP
            • REGEXP_INSTR()
            • REGEXP_LIKE()
            • REGEXP_REPLACE()
            • REGEXP_SUBSTR()
          • aes_decrypt
          • aes_encrypt
          • elt
          • quote
          • right
        • 向量函数
          • 数学计算
          • CLUSTER_CENTERS()
          • COSINE_SIMILARITY()
          • COSINE_DISTANCE()
          • INNER_PRODUCT()
          • L1_NORM()
          • L2_NORM()
          • L2_DISTANCE()
          • NORMALIZE_L2()
          • SUBVECTOR()
          • VECTOR_DIMS()
        • 表函数
          • UNNEST()
        • 窗口函数
          • RANK()
          • ROW_NUMBER()
          • cume_dist
          • percent_rank
        • JSON 函数
          • JSON_EXTRACT()
          • JSON_EXTRACT_FLOAT64()
          • JSON_EXTRACT_STRING()
          • JSON_QUOTE()
          • JSON_ROW()
          • JSON_SET()
          • JSON_UNQUOTE()
          • TRY_JQ()
          • json-arrow
          • JSON_ARRAY()
            • JSON_CONTAINS()
            • JSON_CONTAINS_PATH()
            • JSON_KEYS()
            • JSON_LENGTH()
            • JSON_MERGE_PATCH()
            • JSON_MERGE_PRESERVE()
            • JSON_OBJECT()
            • JSON_OVERLAPS()
            • JSON_PRETTY()
            • JSON_REMOVE()
            • JSON_SCHEMA_VALID()
            • JSON_SCHEMA_VALIDATION_REPORT()
            • JSON_TYPE()
            • JSON_VALID()
            • JSON_VALUE()
          • JSON_KEYS()
          • JSON_LENGTH()
          • JSON_OBJECT()
          • JSON_PRETTY()
          • JSON_SCHEMA_VALID()
          • JSON_SCHEMA_VALIDATION_REPORT()
          • JSON_TYPE()
          • JSON_VALID()
          • JSON_VALUE()
        • 其他函数
          • SAVE_FILE
          • SAMPLE
          • SERIAL_EXTRACT
          • SLEEP
          • STAGE_LIST
          • ONNX_RUN()
          • UUID()
        • 系统运维函数
          • CURRENT_ROLE()
          • CURRENT_USER_NAME()
          • CURRENT_USER()
          • PURGE_LOG()
          • GET_LOCK()
          • RELEASE_LOCK()
          • IS_FREE_LOCK()
          • IS_USED_LOCK()
          • RELEASE_ALL_LOCKS()
          • version
      • 聚合函数
        • AVG
        • APPROX_PERCENTILE
        • BITMAP
        • BIT_AND
        • BIT_OR
        • BIT_XOR
        • COUNT
        • GROUP_CONCAT
        • HLL_ADD_AGG
        • HLL_CARDINALITY
        • HLL_MERGE_AGG
        • MAX
        • MAX_BY
        • MAX_BY_NON_NULL
        • MEDIAN
        • MIN
        • STDDEV_POP
        • SUM
        • VARIANCE
        • VAR_POP
      • 日期时间类
        • CURDATE()
        • CURRENT_TIMESTAMP()
        • DATE()
        • DATE_ADD()
        • DATE_FORMAT()
        • DATE_SUB()
        • DATE_TRUNC()
        • DATEDIFF()
        • DAY()
        • DAYOFYEAR()
        • EXTRACT()
        • HOUR()
        • FROM_UNIXTIME
        • MINUTE()
        • MONTH()
        • NOW()
        • SECOND()
        • STR_TO_DATE()
        • SYSDATE()
        • TIME()
        • TIMEDIFF()
        • TIMESTAMP()
        • TIMESTAMPDIFF()
        • TO_DATE()
        • TO_DAYS()
        • TO_SECONDS()
        • UNIX_TIMESTAMP
        • UTC_TIMESTAMP()
        • WEEK()
        • WEEKDAY()
        • YEAR()
        • addtime
        • curtime
        • dayname
        • get-format
        • maketime
        • monthname
        • quarter
        • subtime
        • time-format
        • timestampadd
        • yearweek
      • 数学类
        • ACOS()
        • ATAN()
        • BIT_COUNT()
        • CEIL()
        • CEILING()
        • COS()
        • COT()
        • CRC32()
        • EXP()
        • FLOOR()
        • LN()
        • LOG()
        • LOG2()
        • LOG10()
        • PI()
        • POWER()
        • ROUND()
        • RAND()
        • SIN()
        • SINH()
        • TAN()
        • atan2
        • degrees
        • radians
        • sign
        • truncate
      • 字符串类
        • BIT_LENGTH()
        • CHAR_LENGTH()
        • CONCAT()
        • CONCAT_WS()
        • EMPTY()
        • ENDSWITH()
        • FIELD()
        • FIND_IN_SET()
        • FORMAT()
        • FROM_BASE64()
        • HEX()
        • INSTR()
        • LCASE()
        • LEFT()
        • LENGTH()
        • LOCATE()
        • LOWER()
        • LPAD()
        • LTRIM()
        • MD5()
        • NAME_CONST()
        • OCT()
        • REPEAT()
        • REVERSE()
        • RPAD()
        • RTRIM()
        • SHA1()/SHA()
        • SHA2()
        • SPACE()
        • SPLIT_PART()
        • STARTSWITH()
        • STRCMP()
        • SUBSTRING()
        • SUBSTRING_INDEX()
        • TO_BASE64()
        • TRIM()
        • UCASE()
        • UNHEX()
        • UPPER()
        • 正则表达式
          • NOT REGEXP
          • REGEXP_INSTR()
          • REGEXP_LIKE()
          • REGEXP_REPLACE()
          • REGEXP_SUBSTR()
        • aes_decrypt
        • aes_encrypt
        • elt
        • quote
        • right
      • 向量类
        • 数学计算
        • CLUSTER_CENTERS()
        • COSINE_SIMILARITY()
        • COSINE_DISTANCE()
        • INNER_PRODUCT()
        • L1_NORM()
        • L2_NORM()
        • L2_DISTANCE()
        • NORMALIZE_L2()
        • SUBVECTOR()
        • VECTOR_DIMS()
      • 表函数
        • UNNEST()
      • 窗口函数
        • RANK()
        • ROW_NUMBER()
        • cume_dist
        • percent_rank
      • JSON 函数
        • JSON_EXTRACT()
        • JSON_EXTRACT_FLOAT64()
        • JSON_EXTRACT_STRING()
        • JSON_QUOTE()
        • JSON_ROW()
        • JSON_SET()
        • JSON_UNQUOTE()
        • TRY_JQ()
        • json-arrow
        • JSON_ARRAY()
          • JSON_CONTAINS()
          • JSON_CONTAINS_PATH()
          • JSON_KEYS()
          • JSON_LENGTH()
          • JSON_MERGE_PATCH()
          • JSON_MERGE_PRESERVE()
          • JSON_OBJECT()
          • JSON_OVERLAPS()
          • JSON_PRETTY()
          • JSON_REMOVE()
          • JSON_SCHEMA_VALID()
          • JSON_SCHEMA_VALIDATION_REPORT()
          • JSON_TYPE()
          • JSON_VALID()
          • JSON_VALUE()
        • JSON_KEYS()
        • JSON_LENGTH()
        • JSON_OBJECT()
        • JSON_PRETTY()
        • JSON_SCHEMA_VALID()
        • JSON_SCHEMA_VALIDATION_REPORT()
        • JSON_TYPE()
        • JSON_VALID()
        • JSON_VALUE()
      • 其他函数
        • SAVE_FILE
        • SAMPLE
        • SERIAL_EXTRACT
        • SLEEP
        • STAGE_LIST
        • ONNX_RUN()
        • UUID()
      • 系统运维函数
        • CURRENT_ROLE()
        • CURRENT_USER_NAME()
        • CURRENT_USER()
        • PURGE_LOG()
        • GET_LOCK()
        • RELEASE_LOCK()
        • IS_FREE_LOCK()
        • IS_USED_LOCK()
        • RELEASE_ALL_LOCKS()
        • version
    • 函数
      • 聚合函数
        • AVG
        • APPROX_PERCENTILE
        • BITMAP
        • BIT_AND
        • BIT_OR
        • BIT_XOR
        • COUNT
        • GROUP_CONCAT
        • HLL_ADD_AGG
        • HLL_CARDINALITY
        • HLL_MERGE_AGG
        • MAX
        • MAX_BY
        • MAX_BY_NON_NULL
        • MEDIAN
        • MIN
        • STDDEV_POP
        • SUM
        • VARIANCE
        • VAR_POP
      • 日期时间
        • CURDATE()
        • CURRENT_TIMESTAMP()
        • DATE()
        • DATE_ADD()
        • DATE_FORMAT()
        • DATE_SUB()
        • DATE_TRUNC()
        • DATEDIFF()
        • DAY()
        • DAYOFYEAR()
        • EXTRACT()
        • HOUR()
        • FROM_UNIXTIME
        • MINUTE()
        • MONTH()
        • NOW()
        • SECOND()
        • STR_TO_DATE()
        • SYSDATE()
        • TIME()
        • TIMEDIFF()
        • TIMESTAMP()
        • TIMESTAMPDIFF()
        • TO_DATE()
        • TO_DAYS()
        • TO_SECONDS()
        • UNIX_TIMESTAMP
        • UTC_TIMESTAMP()
        • WEEK()
        • WEEKDAY()
        • YEAR()
        • addtime
        • curtime
        • dayname
        • get-format
        • maketime
        • monthname
        • quarter
        • subtime
        • time-format
        • timestampadd
        • yearweek
      • Geo 函数
        • H3 Index Functions
        • MBRContains
        • MBRCoveredBy
        • MBRCovers
        • MBRDisjoint
        • MBREquals
        • MBRIntersects
        • MBROverlaps
        • MBRTouches
        • MBRWithin
        • S2 Cell Functions
        • ST_Area()
        • ST_AsGeoJSON
        • ST_AsText()
        • ST_AsWKB()
        • ST_Boundary()
        • ST_Buffer()
        • ST_Centroid()
        • ST_Collect()
        • ST_Contains()
        • ST_ConvexHull()
        • ST_CoveredBy()
        • ST_Covers()
        • ST_Crosses()
        • ST_Difference()
        • ST_Dimension()
        • ST_Disjoint()
        • ST_Distance()
        • ST_Distance_Sphere()
        • ST_EndPoint()
        • ST_Envelope()
        • ST_Equals()
        • ST_ExteriorRing()
        • ST_FrechetDistance()
        • ST_GeoHash
        • ST_GeomCollFromText()
        • ST_GeomCollFromWKB()
        • ST_GeometryN()
        • ST_GeometryType()
        • ST_GeomFromGeoJSON
        • ST_GeomFromText()
        • ST_GeomFromWKB()
        • ST_HausdorffDistance()
        • ST_InteriorRingN()
        • ST_Intersection()
        • ST_Intersects()
        • ST_IsClosed()
        • ST_IsCollection()
        • ST_IsEmpty()
        • ST_IsRing()
        • ST_IsSimple()
        • ST_IsValid()
        • ST_LatFromGeoHash
        • ST_Latitude()
        • ST_Length()
        • ST_LineFromText()
        • ST_LineFromWKB()
        • ST_LineInterpolatePoint()
        • ST_LineInterpolatePoints()
        • ST_LongFromGeoHash
        • ST_Longitude()
        • ST_MakeEnvelope()
        • ST_MLineFromText()
        • ST_MLineFromWKB()
        • ST_MPointFromText()
        • ST_MPointFromWKB()
        • ST_MPolyFromText()
        • ST_MPolyFromWKB()
        • ST_NumGeometries()
        • ST_NumInteriorRings()
        • ST_NumPoints()
        • ST_Overlaps()
        • ST_PointAtDistance()
        • ST_PointFromGeoHash
        • ST_PointFromText()
        • ST_PointFromWKB()
        • ST_PointN()
        • ST_PointOnSurface()
        • ST_PolyFromText()
        • ST_PolyFromWKB()
        • ST_Simplify()
        • ST_SRID()
        • ST_StartPoint()
        • ST_SwapXY()
        • ST_SymDifference()
        • ST_Touches()
        • ST_Union()
        • ST_Validate()
        • ST_Within()
        • ST_X()
        • ST_Y()
      • 数学函数
        • ACOS()
        • ATAN()
        • BIT_COUNT()
        • CEIL()
        • CEILING()
        • COS()
        • COT()
        • CRC32()
        • EXP()
        • FLOOR()
        • LN()
        • LOG()
        • LOG2()
        • LOG10()
        • PI()
        • POWER()
        • ROUND()
        • RAND()
        • SIN()
        • SINH()
        • TAN()
        • atan2
        • degrees
        • radians
        • sign
        • truncate
      • 字符串函数
        • BIT_LENGTH()
        • CHAR_LENGTH()
        • CONCAT()
        • CONCAT_WS()
        • EMPTY()
        • ENDSWITH()
        • FIELD()
        • FIND_IN_SET()
        • FORMAT()
        • FROM_BASE64()
        • HEX()
        • INSTR()
        • LCASE()
        • LEFT()
        • LENGTH()
        • LOCATE()
        • LOWER()
        • LPAD()
        • LTRIM()
        • MD5()
        • NAME_CONST()
        • OCT()
        • REPEAT()
        • REVERSE()
        • RPAD()
        • RTRIM()
        • SHA1()/SHA()
        • SHA2()
        • SPACE()
        • SPLIT_PART()
        • STARTSWITH()
        • STRCMP()
        • SUBSTRING()
        • SUBSTRING_INDEX()
        • TO_BASE64()
        • TRIM()
        • UCASE()
        • UNHEX()
        • UPPER()
        • 正则表达式
          • NOT REGEXP
          • REGEXP_INSTR()
          • REGEXP_LIKE()
          • REGEXP_REPLACE()
          • REGEXP_SUBSTR()
        • aes_decrypt
        • aes_encrypt
        • elt
        • quote
        • right
      • 向量函数
        • 数学计算
        • CLUSTER_CENTERS()
        • COSINE_SIMILARITY()
        • COSINE_DISTANCE()
        • INNER_PRODUCT()
        • L1_NORM()
        • L2_NORM()
        • L2_DISTANCE()
        • NORMALIZE_L2()
        • SUBVECTOR()
        • VECTOR_DIMS()
      • 表函数
        • UNNEST()
      • 窗口函数
        • RANK()
        • ROW_NUMBER()
        • cume_dist
        • percent_rank
      • JSON 函数
        • JSON_EXTRACT()
        • JSON_EXTRACT_FLOAT64()
        • JSON_EXTRACT_STRING()
        • JSON_QUOTE()
        • JSON_ROW()
        • JSON_SET()
        • JSON_UNQUOTE()
        • TRY_JQ()
        • json-arrow
        • JSON_ARRAY()
          • JSON_CONTAINS()
          • JSON_CONTAINS_PATH()
          • JSON_KEYS()
          • JSON_LENGTH()
          • JSON_MERGE_PATCH()
          • JSON_MERGE_PRESERVE()
          • JSON_OBJECT()
          • JSON_OVERLAPS()
          • JSON_PRETTY()
          • JSON_REMOVE()
          • JSON_SCHEMA_VALID()
          • JSON_SCHEMA_VALIDATION_REPORT()
          • JSON_TYPE()
          • JSON_VALID()
          • JSON_VALUE()
        • JSON_KEYS()
        • JSON_LENGTH()
        • JSON_OBJECT()
        • JSON_PRETTY()
        • JSON_SCHEMA_VALID()
        • JSON_SCHEMA_VALIDATION_REPORT()
        • JSON_TYPE()
        • JSON_VALID()
        • JSON_VALUE()
      • 其他函数
        • SAVE_FILE
        • SAMPLE
        • SERIAL_EXTRACT
        • SLEEP
        • STAGE_LIST
        • ONNX_RUN()
        • UUID()
      • 系统运维函数
        • CURRENT_ROLE()
        • CURRENT_USER_NAME()
        • CURRENT_USER()
        • PURGE_LOG()
        • GET_LOCK()
        • RELEASE_LOCK()
        • IS_FREE_LOCK()
        • IS_USED_LOCK()
        • RELEASE_ALL_LOCKS()
        • version
    • 系统配置
      • 单机版通用参数配置
      • 分布式版通用参数配置
    • 系统表目录
    • 权限分类列表
    • 使用限制
      • load-data-support
    • MatrixOne 文件目录结构
    • MatrixOne 工具
      • mo_ctl 分布式工具
      • mo_datax_writer 工具
      • mo_ssb_open 工具
      • mo_tpch_open 工具
      • mo_ts_perf_test 工具
      • mo_service
  • 故障诊断
    • 常用统计数据查询
    • 数据库统计信息
    • 错误码
  • 常见问题解答
    • 部署常见问题
    • SQL 常见问题
  • 版本发布纪要
    • MatrixOne v26.4.2.1 发布说明
    • MatrixOne v26.4.2.0 发布说明
    • MatrixOne v26.4.1.4 发布说明
    • MatrixOne v26.4.1.3 发布说明
    • MatrixOne v26.4.1.2 发布说明
    • MatrixOne v26.4.1.1 发布说明
    • MatrixOne v26.4.1.0 发布说明
    • MatrixOne v26.4.0.0-rc4 发布说明
    • MatrixOne v26.4.0.0-rc3 发布说明
    • MatrixOne v26.4.0.0-rc2 发布说明
    • MatrixOne v26.4.0.0-rc1 发布说明
    • MatrixOne v26.3.0.10 发布报告
    • MatrixOne v26.3.0.11 发布报告
    • MatrixOne v26.3.0.12 发布报告
    • MatrixOne v26.3.0.13 发布报告
    • MatrixOne v26.3.0.14 发布报告
    • MatrixOne v26.3.0.15 发布说明
    • MatrixOne v25.3.0.5 发布报告
    • MatrixOne v25.3.0.6 发布报告
    • MatrixOne v25.3.0.7 发布报告
    • MatrixOne v25.3.0.8 发布报告
    • MatrixOne v25.3.0.9 发布报告
    • MatrixOne v24.2.0.2 发布报告
    • MatrixOne v25.2.0.3 发布报告
    • MatrixOne v25.2.1.0 Release Note
    • MatrixOne v25.2.1.1 发布报告
    • MatrixOne v25.2.2.0 发布报告
    • MatrixOne v25.2.2.1 发布报告
    • MatrixOne v25.2.2.2 发布报告
    • MatrixOne v25.3.0.0 发布声明
    • MatrixOne v25.3.0.1 发布报告
    • MatrixOne v25.3.0.2 发布报告
    • MatrixOne v25.3.0.3 发布报告
    • MatrixOne v25.3.0.4 发布报告
    • MatrixOne v24.1.1.1 发布报告
    • MatrixOne v24.1.1.2 发布报告
    • MatrixOne v24.1.1.3 发布报告
    • MatrixOne v24.1.2.0 发布报告
    • MatrixOne v24.1.2.1 发布报告
    • MatrixOne v24.1.2.2 发布报告
    • MatrixOne v24.1.2.3 发布报告
    • MatrixOne v24.1.2.4 发布报告
    • MatrixOne v24.2.0.0 发布报告
    • MatrixOne v24.2.0.1 发布报告
    • MatrixOne v23.0.7.0 发布报告
    • MatrixOne v23.0.8.0 发布报告
    • MatrixOne v23.1.0.0 发布报告
    • MatrixOne v23.1.0.0-RC1 发布报告
    • MatrixOne v23.1.0.0-RC2 发布报告
    • MatrixOne v23.1.0.1 发布报告
    • MatrixOne v23.1.0.2 发布报告
    • MatrixOne v23.1.1.0 发布报告
    • MatrixOne v22.0.2.0 发布报告
    • MatrixOne v22.0.3.0 发布报告
    • MatrixOne v22.0.4.0 发布报告
    • MatrixOne v22.0.5.0 发布报告
    • MatrixOne v22.0.5.1 发布报告
    • MatrixOne v22.0.6.0 发布报告
    • MatrixOne v21.0.1.0 Release Notes
  • 名词术语表
  • 社区贡献指南
    • 贡献指南
      • 贡献准备
      • 报告 Issue
      • 贡献代码
      • 审核修改
      • 文档贡献
      • 提出设计草案
    • 代码规范
      • 注释规范
      • 提交规范
Back to top

使用 Flink 将 Kafka 数据写入 MatrixOne¶

本章节将介绍如何使用 Flink 将 Kafka 数据写入到 MatrixOne。

前期准备¶

本次实践需要安装部署以下软件环境:

  • 完成单机部署 MatrixOne。

  • 下载安装 lntelliJ IDEA(2022.2.1 or later version)。

  • 根据你的系统环境选择 JDK 8+ version 版本进行下载安装。

  • 下载并安装 Kafka。

  • 下载并安装 Flink,最低支持版本为 1.11。

  • 下载并安装 MySQL Client。

操作步骤¶

步骤一:启动 Kafka 服务¶

Kafka 集群协调和元数据管理可以通过 KRaft 或 ZooKeeper 来实现。在这里,我们将使用 Kafka 3.5.0 版本,无需依赖独立的 ZooKeeper 软件,而是使用 Kafka 自带的 KRaft 来进行元数据管理。请按照以下步骤配置配置文件,该文件位于 Kafka 软件根目录下的 config/kraft/server.properties。

配置文件内容如下:

# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements.  See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License.  You may obtain a copy of the License at
#
#    http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

#
# This configuration file is intended for use in KRaft mode, where
# Apache ZooKeeper is not present.  See config/kraft/README.md for details.
#

############################# Server Basics #############################

# The role of this server. Setting this puts us in KRaft mode
process.roles=broker,controller

# The node id associated with this instance's roles
node.id=1

# The connect string for the controller quorum
controller.quorum.voters=1@xx.xx.xx.xx:9093

############################# Socket Server Settings #############################

# The address the socket server listens on.
# Combined nodes (i.e. those with `process.roles=broker,controller`) must list the controller listener here at a minimum.
# If the broker listener is not defined, the default listener will use a host name that is equal to the value of java.net.InetAddress.getCanonicalHostName(),
# with PLAINTEXT listener name, and port 9092.
#   FORMAT:
#     listeners = listener_name://host_name:port
#   EXAMPLE:
#     listeners = PLAINTEXT://your.host.name:9092
#listeners=PLAINTEXT://:9092,CONTROLLER://:9093
listeners=PLAINTEXT://xx.xx.xx.xx:9092,CONTROLLER://xx.xx.xx.xx:9093

# Name of listener used for communication between brokers.
inter.broker.listener.name=PLAINTEXT

# Listener name, hostname and port the broker will advertise to clients.
# If not set, it uses the value for "listeners".
#advertised.listeners=PLAINTEXT://localhost:9092

# A comma-separated list of the names of the listeners used by the controller.
# If no explicit mapping set in `listener.security.protocol.map`, default will be using PLAINTEXT protocol
# This is required if running in KRaft mode.
controller.listener.names=CONTROLLER

# Maps listener names to security protocols, the default is for them to be the same. See the config documentation for more details
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL

# The number of threads that the server uses for receiving requests from the network and sending responses to the network
num.network.threads=3

# The number of threads that the server uses for processing requests, which may include disk I/O
num.io.threads=8

# The send buffer (SO_SNDBUF) used by the socket server
socket.send.buffer.bytes=102400

# The receive buffer (SO_RCVBUF) used by the socket server
socket.receive.buffer.bytes=102400

# The maximum size of a request that the socket server will accept (protection against OOM)
socket.request.max.bytes=104857600


############################# Log Basics #############################

# A comma separated list of directories under which to store log files
log.dirs=/home/software/kafka_2.13-3.5.0/kraft-combined-logs

# The default number of log partitions per topic. More partitions allow greater
# parallelism for consumption, but this will also result in more files across
# the brokers.
num.partitions=1

# The number of threads per data directory to be used for log recovery at startup and flushing at shutdown.
# This value is recommended to be increased for installations with data dirs located in RAID array.
num.recovery.threads.per.data.dir=1

############################# Internal Topic Settings  #############################
# The replication factor for the group metadata internal topics "__consumer_offsets" and "__transaction_state"
# For anything other than development testing, a value greater than 1 is recommended to ensure availability such as 3.
offsets.topic.replication.factor=1
transaction.state.log.replication.factor=1
transaction.state.log.min.isr=1

############################# Log Flush Policy #############################

# Messages are immediately written to the filesystem but by default we only fsync() to sync
# the OS cache lazily. The following configurations control the flush of data to disk.
# There are a few important trade-offs here:
#    1. Durability: Unflushed data may be lost if you are not using replication.
#    2. Latency: Very large flush intervals may lead to latency spikes when the flush does occur as there will be a lot of data to flush.
#    3. Throughput: The flush is generally the most expensive operation, and a small flush interval may lead to excessive seeks.
# The settings below allow one to configure the flush policy to flush data after a period of time or
# every N messages (or both). This can be done globally and overridden on a per-topic basis.

# The number of messages to accept before forcing a flush of data to disk
#log.flush.interval.messages=10000

# The maximum amount of time a message can sit in a log before we force a flush
#log.flush.interval.ms=1000

############################# Log Retention Policy #############################

# The following configurations control the disposal of log segments. The policy can
# be set to delete segments after a period of time, or after a given size has accumulated.
# A segment will be deleted whenever *either* of these criteria are met. Deletion always happens
# from the end of the log.

# The minimum age of a log file to be eligible for deletion due to age
log.retention.hours=72

# A size-based retention policy for logs. Segments are pruned from the log unless the remaining
# segments drop below log.retention.bytes. Functions independently of log.retention.hours.
#log.retention.bytes=1073741824

# The maximum size of a log segment file. When this size is reached a new log segment will be created.
log.segment.bytes=1073741824

# The interval at which log segments are checked to see if they can be deleted according
# to the retention policies
log.retention.check.interval.ms=300000

文件配置完成后,执行如下命令,启动 Kafka 服务:

#生成集群 ID
$ KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
#设置日志目录格式
$ bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties
#启动 Kafka 服务
$ bin/kafka-server-start.sh config/kraft/server.properties

步骤二:创建 Kafka 主题¶

为了使 Flink 能够从中读取数据并写入到 MatrixOne,我们需要首先创建一个名为 “matrixone” 的 Kafka 主题。在下面的命令中,使用 --bootstrap-server 参数指定 Kafka 服务的监听地址为 xx.xx.xx.xx:9092:

$ bin/kafka-topics.sh --create --topic matrixone --bootstrap-server xx.xx.xx.xx:9092

步骤三:读取 MatrixOne 数据¶

在连接到 MatrixOne 数据库之后,需要执行以下操作以创建所需的数据库和数据表:

  1. 在 MatrixOne 中创建数据库和数据表,并导入数据:

    CREATE TABLE `users` (
    `id` INT DEFAULT NULL,
    `name` VARCHAR(255) DEFAULT NULL,
    `age` INT DEFAULT NULL
    )
    
  2. 在 IDEA 集成开发环境中编写代码:

    在 IDEA 中,创建两个类:User.java 和 Kafka2Mo.java。这些类用于使用 Flink 从 Kafka 读取数据,并将数据写入 MatrixOne 数据库中。

package com.matrixone.flink.demo.entity;

public class User {

    private int id;
    private String name;
    private int age;

    public int getId() {
        return id;
    }

    public void setId(int id) {
        this.id = id;
    }

    public String getName() {
        return name;
    }

    public void setName(String name) {
        this.name = name;
    }

    public int getAge() {
        return age;
    }

    public void setAge(int age) {
        this.age = age;
    }
}
package com.matrixone.flink.demo;

import com.alibaba.fastjson2.JSON;
import com.matrixone.flink.demo.entity.User;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.serialization.AbstractDeserializationSchema;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.apache.flink.connector.jdbc.JdbcStatementBuilder;
import org.apache.flink.connector.jdbc.internal.options.JdbcConnectorOptions;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.kafka.clients.consumer.OffsetResetStrategy;

import java.nio.charset.StandardCharsets;

/**
 * @author MatrixOne
 * @desc
 */
public class Kafka2Mo {

    private static String srcServer = "xx.xx.xx.xx:9092";
    private static String srcTopic = "matrixone";
    private static String consumerGroup = "matrixone_group";

    private static String destHost = "xx.xx.xx.xx";
    private static Integer destPort = 6001;
    private static String destUserName = "root";
    private static String destPassword = "111";
    private static String destDataBase = "test";
    private static String destTable = "person";

    public static void main(String[] args) throws Exception {

        //初始化环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        //设置并行度
        env.setParallelism(1);

        //设置 kafka source 信息
        KafkaSource<User> source = KafkaSource.<User>builder()
                //Kafka 服务
                .setBootstrapServers(srcServer)
                //消息主题
                .setTopics(srcTopic)
                //消费组
                .setGroupId(consumerGroup)
                //偏移量 当没有提交偏移量则从最开始开始消费
                .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.LATEST))
                //自定义解析消息内容
                .setValueOnlyDeserializer(new AbstractDeserializationSchema<User>() {
                    @Override
                    public User deserialize(byte[] message) {
                        return JSON.parseObject(new String(message, StandardCharsets.UTF_8), User.class);
                    }
                })
                .build();
        DataStreamSource<User> kafkaSource = env.fromSource(source, WatermarkStrategy.noWatermarks(), "kafka_maxtixone");
        //kafkaSource.print();

        //设置 matrixone sink 信息
        kafkaSource.addSink(JdbcSink.sink(
                "insert into users (id,name,age) values(?,?,?)",
                (JdbcStatementBuilder<User>) (preparedStatement, user) -> {
                    preparedStatement.setInt(1, user.getId());
                    preparedStatement.setString(2, user.getName());
                    preparedStatement.setInt(3, user.getAge());
                },
                JdbcExecutionOptions.builder()
                        //默认值 5000
                        .withBatchSize(1000)
                        //默认值为 0
                        .withBatchIntervalMs(200)
                        //最大尝试次数
                        .withMaxRetries(5)
                        .build(),
                JdbcConnectorOptions.builder()
                        .setDBUrl("jdbc:mysql://"+destHost+":"+destPort+"/"+destDataBase)
                        .setUsername(destUserName)
                        .setPassword(destPassword)
                        .setDriverName("com.mysql.cj.jdbc.Driver")
                        .setTableName(destTable)
                        .build()
        ));
        env.execute();
    }
}

代码编写完成后,你可以运行 Flink 任务,即在 IDEA 中选择 Kafka2Mo.java 文件,然后执行 Kafka2Mo.Main()。

步骤四:生成数据¶

使用 Kafka 提供的命令行生产者工具,您可以向 Kafka 的 “matrixone” 主题中添加数据。在下面的命令中,使用 --topic 参数指定要添加到的主题,而 --bootstrap-server 参数指定了 Kafka 服务的监听地址。

bin/kafka-console-producer.sh --topic matrixone --bootstrap-server xx.xx.xx.xx:9092

执行上述命令后,您将在控制台上等待输入消息内容。只需直接输入消息值 (value),每行表示一条消息(以换行符分隔),如下所示:

{"id": 10, "name": "xiaowang", "age": 22}
{"id": 20, "name": "xiaozhang", "age": 24}
{"id": 30, "name": "xiaogao", "age": 18}
{"id": 40, "name": "xiaowu", "age": 20}
{"id": 50, "name": "xiaoli", "age": 42}

步骤五:查看执行结果¶

在 MatrixOne 中执行如下 SQL 查询结果:

mysql> select * from test.users;
+------+-----------+------+
| id   | name      | age  |
+------+-----------+------+
|   10 | xiaowang  |   22 |
|   20 | xiaozhang |   24 |
|   30 | xiaogao   |   18 |
|   40 | xiaowu    |   20 |
|   50 | xiaoli    |   42 |
+------+-----------+------+
5 rows in set (0.01 sec)
Next
使用 DolphinScheduler 连接 MatrixOne
Previous
使用 Flink 将 TiDB 数据写入 MatrixOne
Copyright © 2026, MatrixOrigin
Made with Sphinx and @pradyunsg's Furo
On this page
  • 使用 Flink 将 Kafka 数据写入 MatrixOne
    • 前期准备
    • 操作步骤
      • 步骤一:启动 Kafka 服务
      • 步骤二:创建 Kafka 主题
      • 步骤三:读取 MatrixOne 数据
      • 步骤四:生成数据
      • 步骤五:查看执行结果