ClickHouse 数据湖测试数据生成指南:使用 Paimon Java 客户端构造 Paimon 格式目录 📅 发布时间:2026/9/21 2:18:19 👁 浏览次数: 数据库OLAP列式数据库大数据实时分析数据分析【免费下载链接】ClickHouseClickHouse® is a real-time analytics database management system项目地址https://gitcode.com/GitHub_Trending/cli/ClickHouse点击查看免费下载导读本文讲解如何从零构建一个符合 Apache Paimon 布局规范的数据库目录含schema/、snapshot/、manifest/与分区数据文件用于 ClickHouse 数据湖场景的功能测试与验证。你将掌握基于 Maven JDK 17 编写 Paimon Java 数据生成器、覆盖全部基础类型含可空与非空变体、ARRAY、MAP的完整流程并理解生成结果在 ClickHousepaimonS3表函数与PaimonS3存储引擎中如何被消费和校验。Paimon 测试数据在 ClickHouse 仓库中的角色在 ClickHouse 仓库中tests/queries/0_stateless/data_minio/目录下存放着一批面向对象存储MinIO/AWS S3的静态测试数据集其中paimon_all_types/专门用于验证 ClickHouse 对 Paimon 表格式的全类型读取能力。该目录不是手工拼装的假目录而是由 Paimon 官方 Java 客户端按标准写入流程生成的真实 Paimon 表数据目录根下存在schema/schema-0表结构定义、snapshot/snapshot-1与EARLIEST/LATEST指针快照状态、manifest/manifest-list-*.avro与manifest-*.avro清单列表与数据清单数据按 20 个分区键逐层嵌套如f_booleantrue/f_string中文String0/f_date19358/...最终落在bucket-0/data-*.parquet数据文件中。这份数据集的生成方式正是tests/queries/0_stateless/data_minio/paimon_all_types/README.md所记录的内容。它同时服务于多条无状态测试例如 03546_paimon_all_supported_type.sql 通过paimonS3(s3_conn, filenamepaimon_all_types)表函数校验 40 列的读取结果03775_paimon_incremental_read.sql 则基于同一数据集验证增量读取能力。因此掌握如何生成这类目录是扩展 ClickHouse 数据湖测试矩阵的基础技能。一、环境前置要求README 明确给出了数据生成器开发环境的两项硬性要求组件版本README 中实测记录说明Apache Maven3.9.9commit8e8579a9e76f7d015ee5ec7bfcdc97d260186937构建与执行 Java 工程JDKjava 17.0.12 2024-07-16 LTSPaimon/Flink 生态依赖 Java 17此外由于数据生成依赖 Flink Table 生态paimon-flink-common、flink-table-*构建过程会拉取大量依赖请确保网络可访问 Maven 中央仓库。二、创建 Maven 工程与 pom.xml2.1 初始化工程使用你习惯的方式创建一个 Maven 工程例如mvn archetype:generate或 IDE 向导工程结构最终如下README 中给出的最终项目树paimon-example/ ├── pom.xml └── src └── main └── java └── org └── apache └── paimon └── service └── example └── DataGenerator.java即 9 个目录、2 个文件target构建产物除外。2.2 完整 pom.xml创建pom.xml将mainClass替换为你的实际主类名?xml version1.0 encodingUTF-8? project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion version1.1.1/version groupIdorg.apache.paimon/groupId artifactIdpaimon-example/artifactId properties project.build.sourceEncodingUTF-8/project.build.sourceEncoding hadoop.version2.8.5/hadoop.version log4j.version2.17.1/log4j.version /properties dependencies dependency groupIdorg.apache.paimon/groupId artifactIdpaimon-common/artifactId version${project.version}/version /dependency dependency groupIdorg.apache.paimon/groupId artifactIdpaimon-core/artifactId version${project.version}/version /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-common/artifactId version${hadoop.version}/version scoperuntime/scope exclusions exclusion groupIdlog4j/groupId artifactIdlog4j/artifactId /exclusion exclusion groupIdorg.slf4j/groupId artifactIdslf4j-log4j12/artifactId /exclusion /exclusions /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-hdfs-client/artifactId version${hadoop.version}/version scoperuntime/scope exclusions exclusion groupIdlog4j/groupId artifactIdlog4j/artifactId /exclusion exclusion groupIdorg.slf4j/groupId artifactIdslf4j-log4j12/artifactId /exclusion /exclusions /dependency dependency groupIdorg.apache.paimon/groupId artifactIdpaimon-format/artifactId version${project.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-common/artifactId version1.20.1/version scopecompile/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version1.20.1/version scopecompile/scope /dependency dependency groupIdorg.apache.paimon/groupId artifactIdpaimon-flink-common/artifactId version1.1-SNAPSHOT/version scopecompile/scope /dependency dependency groupIdorg.apache.commons/groupId artifactIdcommons-compress/artifactId version1.24.0/version /dependency !-- https://mvnrepository.com/artifact/org.apache.flink/flink-connector-base -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-base/artifactId version1.20.1/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-runtime/artifactId version1.20.1/version scopecompile/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridge/artifactId version1.20.1/version scopecompile/scope /dependency !-- https://mvnrepository.com/artifact/org.apache.flink/flink-table-planner -- dependency groupIdorg.apache.flink/groupId artifactIdflink-table-planner_2.12/artifactId version1.20.1/version scopecompile/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-clients/artifactId version1.20.1/version scoperuntime/scope /dependency dependency groupIdcommons-io/groupId artifactIdcommons-io/artifactId version2.11.0/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-java/artifactId version1.20.1/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version1.20.1/version /dependency /dependencies build plugins plugin groupIdorg.codehaus.mojo/groupId artifactIdexec-maven-plugin/artifactId version3.0.0/version executions execution goals goaljava/goal /goals /execution /executions configuration addResourcesToClasspathtrue/addResourcesToClasspath mainClassorg.apache.paimon.service.example.DataGenerator/mainClass /configuration /plugin /plugins /build /project几点版本说明来自该 pom 的实际约束paimon-*系列依赖paimon-common/paimon-core/paimon-format使用${project.version}即继承工程自身的version1.1.1/versionpaimon-flink-common使用1.1-SNAPSHOT快照版本需要仓库已发布对应快照Flink 侧统一使用1.20.1且flink-table-planner_2.12与flink-streaming-java同时出现两次README 原始 pom 即如此Scala 2.12 与 Java API 均被引用Hadoop 以runtime作用域引入并显式排除log4j与slf4j-log4j12避免日志门面冲突。三、编写 DataGenerator 数据生成类在src/main/java/org/apache/paimon/service/example/DataGenerator.java中编写主类并将代码中的rootPath替换为目标目录路径README 示例中为/tmp/warehouse。3.1 完整源码package org.apache.paimon.service.example; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; import org.apache.paimon.catalog.Catalog; import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.catalog.CatalogFactory; import org.apache.paimon.catalog.Identifier; import org.apache.paimon.data.*; import org.apache.paimon.disk.IOManagerImpl; import org.apache.paimon.flink.FlinkCatalog; import org.apache.paimon.fs.Path; import org.apache.paimon.options.Options; import org.apache.paimon.schema.Schema; import org.apache.paimon.table.Table; import org.apache.paimon.table.sink.*; import org.apache.paimon.types.*; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.math.BigDecimal; import java.time.LocalDate; import java.time.LocalDateTime; import java.time.LocalTime; import java.util.HashMap; import java.util.List; import java.util.Map; public class DataGenerator { private static final Logger LOG LoggerFactory.getLogger(DataGenerator.class); public static Catalog createFilesystemCatalog(String path) { CatalogContext context CatalogContext.create(new Path(path)); return CatalogFactory.createCatalog(context); } private static Catalog getCatalog(String rootPath) { try { return org.apache.paimon.catalog.CatalogFactory.createCatalog(CatalogContext.create(new Path(rootPath))); } catch (Exception e) { throw new RuntimeException(Create Catalog failed, e); } } private static InternalRow createRow(int id, RowType rowType) { GenericRow row new GenericRow(rowType.getFieldCount()); ListDataType fieldTypes rowType.getFieldTypes(); ListString fieldNames rowType.getFieldNames(); for (int i 0; i fieldNames.size(); i) { String fieldName fieldNames.get(i); DataType fieldType fieldTypes.get(i); switch (fieldName) { case f_boolean: case f_boolean_nn: row.setField(i, id % 2 0); break; case f_char: case f_char_nn: row.setField(i, BinaryString.fromString(String.valueOf((char)(A id%26)))); break; case f_varchar: case f_varchar_nn: row.setField(i, BinaryString.fromString(String.valueOf((char)(a id%26)))); break; case f_string: case f_string_nn: row.setField(i, BinaryString.fromString(中文String id)); break; case f_binary: case f_varbinary: case f_bytes: case f_binary_nn: case f_varbinary_nn: case f_bytes_nn: row.setField(i, new byte[]{(byte)id}); break; case f_decimal: case f_decimal_nn: row.setField(i, Decimal.fromBigDecimal( new BigDecimal(id).setScale(1), ((DecimalType) fieldType).getPrecision(), ((DecimalType) fieldType).getScale() )); break; case f_decimal2: case f_decimal2_nn: row.setField(i, Decimal.fromBigDecimal( new BigDecimal(id * 10L).setScale(1), ((DecimalType) fieldType).getPrecision(), ((DecimalType) fieldType).getScale() )); break; case f_decimal3: case f_decimal3_nn: row.setField(i, Decimal.fromBigDecimal( new BigDecimal(id * 100L).setScale(1), ((DecimalType) fieldType).getPrecision(), ((DecimalType) fieldType).getScale() )); break; case f_tinyint: case f_tinyint_nn: row.setField(i, (byte)id); break; case f_smallint: case f_smallint_nn: row.setField(i, (short)id); break; case f_int: case f_int_nn: row.setField(i, id); break; case f_bigint: case f_bigint_nn: row.setField(i, (long)id * 1000); break; case f_float: case f_float_nn: row.setField(i, (float)id 0.1f); break; case f_double: case f_double_nn: row.setField(i, (double)id 0.01); break; case f_date: case f_date_nn: row.setField(i, (int)(LocalDate.of(2023, 1, Math.max(1, id % 31)).toEpochDay())); break; case f_time: case f_time_nn: LocalTime time LocalTime.of(id % 24, id % 60, id % 60); row.setField(i, time.toSecondOfDay() * 1000); break; case f_timestamp: case f_timestamp2: case f_timestamp3: case f_timestamp_nn: case f_timestamp2_nn: case f_timestamp3_nn: LocalDateTime timestamp LocalDateTime.of(2025, 1, id % 31 1, id % 24, id % 60, id % 60, id * 1000 * 1000); row.setField(i, Timestamp.fromLocalDateTime(timestamp)); break; case f_array: GenericArray arrayData new GenericArray(new int[]{id, id*2, id*3}); row.setField(i, arrayData); break; case f_map: MapBinaryString, BinaryString data new HashMap(); data.put(BinaryString.fromString(Integer.toString(id)), BinaryString.fromString(Integer.toString(id))); data.put(BinaryString.fromString(Integer.toString(id*2)), BinaryString.fromString(Integer.toString(id*2))); data.put(BinaryString.fromString(Integer.toString(id*3)), BinaryString.fromString(Integer.toString(id*3))); GenericMap mapData new GenericMap(data); row.setField(i, mapData); break; default: throw new RuntimeException(unknown column name: fieldName); } if ((!fieldName.endsWith(_nn) !fieldType.is(DataTypeRoot.ARRAY) !fieldType.is(DataTypeRoot.MAP)) id % 2 0) { row.setField(i, null); } } return row; } public static void generateTestCase1(String rootPath) throws Exception { { /// create table Schema.Builder schemaBuilder Schema.newBuilder(); schemaBuilder.column(f_boolean, DataTypes.BOOLEAN()); schemaBuilder.column(f_char, DataTypes.CHAR(1)); schemaBuilder.column(f_varchar, DataTypes.VARCHAR(1)); schemaBuilder.column(f_string, DataTypes.STRING()); schemaBuilder.column(f_binary, DataTypes.BINARY(1)); schemaBuilder.column(f_varbinary, DataTypes.VARBINARY(1)); schemaBuilder.column(f_bytes, DataTypes.BYTES()); schemaBuilder.column(f_decimal, DataTypes.DECIMAL(9, 1)); schemaBuilder.column(f_decimal2, DataTypes.DECIMAL(18, 1)); schemaBuilder.column(f_decimal3, DataTypes.DECIMAL(38, 1)); schemaBuilder.column(f_tinyint, DataTypes.TINYINT()); schemaBuilder.column(f_smallint, DataTypes.SMALLINT()); schemaBuilder.column(f_int, DataTypes.INT()); schemaBuilder.column(f_bigint, DataTypes.BIGINT()); schemaBuilder.column(f_float, DataTypes.FLOAT()); schemaBuilder.column(f_double, DataTypes.DOUBLE()); schemaBuilder.column(f_date, DataTypes.DATE()); schemaBuilder.column(f_time, DataTypes.TIME()); schemaBuilder.column(f_timestamp, DataTypes.TIMESTAMP(3)); schemaBuilder.column(f_timestamp2, DataTypes.TIMESTAMP(1)); schemaBuilder.column(f_timestamp3, DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE(3)); schemaBuilder.column(f_boolean_nn, DataTypes.BOOLEAN().notNull()); schemaBuilder.column(f_char_nn, DataTypes.CHAR(1).notNull()); schemaBuilder.column(f_varchar_nn, DataTypes.VARCHAR(1).notNull()); schemaBuilder.column(f_string_nn, DataTypes.STRING().notNull()); schemaBuilder.column(f_binary_nn, DataTypes.BINARY(1).notNull()); schemaBuilder.column(f_varbinary_nn, DataTypes.VARBINARY(1).notNull()); schemaBuilder.column(f_bytes_nn, DataTypes.BYTES().notNull()); schemaBuilder.column(f_decimal_nn, DataTypes.DECIMAL(9, 1).notNull()); schemaBuilder.column(f_decimal2_nn, DataTypes.DECIMAL(18, 1).notNull()); schemaBuilder.column(f_decimal3_nn, DataTypes.DECIMAL(38, 1).notNull()); schemaBuilder.column(f_tinyint_nn, DataTypes.TINYINT().notNull()); schemaBuilder.column(f_smallint_nn, DataTypes.SMALLINT().notNull()); schemaBuilder.column(f_int_nn, DataTypes.INT().notNull()); schemaBuilder.column(f_bigint_nn, DataTypes.BIGINT().notNull()); schemaBuilder.column(f_float_nn, DataTypes.FLOAT().notNull()); schemaBuilder.column(f_double_nn, DataTypes.DOUBLE().notNull()); schemaBuilder.column(f_date_nn, DataTypes.DATE().notNull()); schemaBuilder.column(f_time_nn, DataTypes.TIME().notNull()); schemaBuilder.column(f_timestamp_nn, DataTypes.TIMESTAMP(3).notNull()); schemaBuilder.column(f_timestamp2_nn, DataTypes.TIMESTAMP(1).notNull()); schemaBuilder.column(f_timestamp3_nn, DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE(3).notNull()); schemaBuilder.column(f_array, DataTypes.ARRAY(DataTypes.INT().notNull()).notNull()); schemaBuilder.column(f_map, DataTypes.MAP(DataTypes.STRING().notNull(), DataTypes.STRING().notNull()).notNull()); schemaBuilder.partitionKeys( f_boolean, f_string, f_date, f_boolean_nn, f_char_nn, f_varchar_nn, f_string_nn, f_decimal_nn, f_decimal2_nn, f_decimal3_nn, f_tinyint_nn, f_smallint_nn, f_int_nn, f_bigint_nn, f_float_nn, f_double_nn, f_date_nn, f_time_nn, f_timestamp_nn, f_timestamp2_nn); Schema schema schemaBuilder.build(); Identifier identifier Identifier.create(tests, cases2); try { Catalog catalog createFilesystemCatalog(rootPath); catalog.createDatabase(tests, true); catalog.createTable(identifier, schema, false); } catch (Catalog.TableAlreadyExistException e) { // do something } catch (Catalog.DatabaseNotExistException e) { // do something } catch (Catalog.DatabaseAlreadyExistException e) { throw new RuntimeException(e); } Catalog catalog getCatalog(rootPath); Identifier tableId Identifier.create(tests, cases2); Table table catalog.getTable(tableId); BatchWriteBuilder writeBuilder table.newBatchWriteBuilder(); TableWriteImpl writer (TableWriteImpl) writeBuilder.newWrite() .withIOManager(new IOManagerImpl(rootPath)); for (int i 0; i 10; i) { InternalRow row createRow(i, table.rowType()); writer.write(row); } ListCommitMessage messages writer.prepareCommit(); BatchTableCommit commit writeBuilder.newCommit(); commit.commit(messages); } } public static void main(String[] args) throws Exception { generateTestCase1(/tmp/warehouse); } }3.2 关键实现点剖析Catalog 创建Filesystem CatalogcreateFilesystemCatalog与getCatalog都通过CatalogFactory.createCatalog(CatalogContext.create(new Path(rootPath)))以本地文件系统路径构建 Catalog这是 Paimon Catalog 即目录 设计的体现——最终rootPath下会同时包含 Catalog 元数据与表数据。两者的差异在于异常处理方式前者直接返回后者将异常包装为RuntimeException抛出。行数据构造createRow根据表的RowType反射式地为每个字段赋值覆盖 20 类字段族含_nn非空后缀变体。值得注意的取值设计f_boolean/f_boolean_nn取id % 2 0使布尔分区呈现true与false两类字符类f_char/f_varchar由A id % 26推导保证 10 行数据生成不同的分区键字符串采用中文String id刻意引入非 ASCII 内容用于验证多字节字符在分区目录名与 Parquet 数据中的往返一致性仓库中的生成结果目录确实出现了f_string中文String0这样的中文目录名时间类中f_date使用LocalDate.of(2023, 1, ...)计算 epoch 天数对应仓库数据目录中19358这样的整数日期f_timestamp*使用 2025 年 1 月日期并带纳秒级偏移以覆盖 TIMESTAMP 精度与本地时区TIMESTAMP_WITH_LOCAL_TIME_ZONE语义。可空性注入赋值循环末尾有一个关键逻辑——对非_nn结尾、且非 ARRAY/MAP的字段在id % 2 0时显式写入null。这意味着 10 行数据中偶数行0、2、4、6、8的所有可空基础类型字段均为 NULL用于验证 ClickHouse 端 NULL 传播与类型推导的正确性而_nn变体始终非空便于对照。ARRAY/MAP 列则始终保持非空规避复合类型的 NULL 编码差异。建表与分区schemaBuilder.partitionKeys(...)指定了 20 个分区键——包含 3 个可空列f_boolean、f_string、f_date与 17 个_nn非空列。将可空列纳入分区键能系统性检验分区裁剪 可空值的交互仓库生成目录中可空分区列存在__DEFAULT_PARTITION__占位正是可空分区列的默认分区名。写入与提交写入流程完整走 Paimon 标准 sink 链路table.newBatchWriteBuilder()→newWrite().withIOManager(new IOManagerImpl(rootPath))→ 逐行writer.write(row)→prepareCommit()产出CommitMessage列表 →newCommit().commit(messages)。数据以bucket-0的固定桶写入data-*.parquet文件并同步生成manifest/清单与snapshot/snapshot-1快照——这正是 ClickHouse 读取 Paimon 表所依赖的三层元数据。四、构建与运行4.1 构建工程mvn install -DskipTests -Dcheckstyle.skiptrue -Dspotless.check.skiptrue -Drat.skiptrue -Denforcer.skip该命令跳过测试及代码风格类插件checkstyle、spotless、rat、enforcer仅完成编译与安装适合本地快速生成数据。4.2 运行数据生成器mvn exec:java -Dexec.mainClassorg.apache.paimon.service.example.Example -Dcheckstyle.skiptrue通过exec-maven-plugin的javagoal 直接运行主类README 中-Dexec.mainClass指向示例类名org.apache.paimon.service.example.Example实际运行时请替换为你的主类org.apache.paimon.service.example.DataGenerator。运行成功后/tmp/warehouse或你配置的rootPath下即生成完整的 Paimon 表目录。4.3 生成目录结构解读README 记录的最终项目树如下target已被tree -I target排除➜ paimon-example git:(main) ✗ tree -I target . . ├── pom.xml └── src └── main └── java └── org └── apache └── paimon └── service └── example └── DataGenerator.java 9 directories, 2 files而运行后产出的Paimon 数据目录非工程目录则可以对照仓库中已生成好的 paimon_all_types 数据集来理解其布局。从该目录tests/queries/0_stateless/data_minio/paimon_all_types/可以确认schema/存放schema-0表结构文件Avro 序列化snapshot/包含snapshot-1快照文件与EARLIEST、LATEST两个指针文件指向当前可用快照manifest/包含manifest-list-*.avro清单列表与manifest-*.avro数据文件清单ClickHouse 通过它们定位 Parquet 数据文件分区目录按分区键值逐层嵌套例如f_booleantrue/f_string中文String0/f_date19358/f_boolean_nntrue/f_char_nnA/...最深处的bucket-0/下才是data-uuid-0.parquet数据文件。仓库中生成的paimon_all_types数据与 README 中generateTestCase1的 schema 定义一一对应20 个分区键、41 个数据列、10 行数据分别落入 10 个不同的分区组合。五、在 ClickHouse 中消费与验证生成的 Paimon 目录生成 Paimon 目录的目的是供 ClickHouse 数据湖读取链路进行测试。仓库中对应的消费入口有两类均可直接对标5.1 表函数 paimonS3一次性读取03546_paimon_all_supported_type.sql 展示了最直接的验证方式-- Tags: no-fasttest依赖 AWS/MinIO SET enable_time_time64_type 1, session_timezone UTC; DESC paimonS3(s3_conn, filename paimon_all_types); SELECT f_boolean, f_char, ..., toTimeZone(f_timestamp3, Asia/Shanghai), ..., f_array, f_map FROM paimonS3(s3_conn, filename paimon_all_types) ORDER BY f_int_nn; SELECT count(1) FROM paimonS3(s3_conn, filename paimon_all_types);该测试验证了类型推导完整DESC输出 40 列、全类型数据可读、TIMESTAMP 时区转换toTimeZone(..., Asia/Shanghai)正确、ORDER BY f_int_nn可对分区列排序、总行数稳定应为 10。5.2 存储引擎 PaimonS3持久化 增量读取03775_paimon_incremental_read.sql 则基于同一数据集测试增量读取语义SET enable_time_time64_type 1, session_timezone UTC, allow_experimental_paimon_storage_engine 1; CREATE TABLE paimon_inc_read ENGINE PaimonS3(s3_conn, filename paimon_all_types) SETTINGS paimon_incremental_read 1, paimon_keeper_path /clickhouse/tables/{database}/paimon_inc_read, paimon_replica_name {replica}; -- First run: 读到 latest snapshot 的增量数据 SELECT count() FROM paimon_inc_read; -- Second run: 无新 snapshot应返回 0 SELECT count() FROM paimon_inc_read;这里paimon_keeper_path、paimon_replica_name用于在多副本间记录已消费的 snapshot 位置实现仅读取新增快照的增量语义。5.3 ClickHouse 侧读取链路的源码印证ClickHouse 的 Paimon 读取实现集中在src/Storages/ObjectStorage/DataLakes/Paimon/目录可从源码结构印证上述测试为何能工作PaimonClient.h 定义了PaimonSnapshot、PaimonManifestFileMeta、PaimonManifestEntry::DataFileMeta与PaimonTableClient其中PaimonTableClient提供getLatestTableSnapshotInfo()、getDataManifest()、getManifestMeta()等方法对应读 snapshot → 读 manifest-list → 读 manifest → 定位 data 文件的读取链路Constant.h 集中定义了读取时使用的 Paimon 元数据列名常量如PAIMON_SNAPSHOT_DIRsnapshot、PAIMON_MANIFEST_DIRmanifest、PAIMON_DEFAULT_PARTITION_NAMEpartition.default-name对应生成目录中可空分区列的__DEFAULT_PARTITION__占位、以及 manifest 中各字段_KIND、_PARTITION、_BUCKET、_FILE等的列名映射。从这些代码可以推断ClickHouse 并不直接解析 Paimon 的 Parquet 分区目录来做谓词推导而是优先读取snapshot/manifest元数据其中包含_MAX_VALUES/_MIN_VALUES/_NULL_COUNTS等统计信息再据此裁剪数据文件——这也是为什么测试数据集必须由 Paimon 官方客户端生成而非手工拼装。六、扩展其他 Paimon 测试数据集的生成思路仓库tests/queries/0_stateless/data_minio/下还提供了一批同类 Paimon 数据集生成方式与paimon_all_types类似都是写 Java 生成器 → 运行 → 产物提交到仓库数据集目录覆盖点对应测试paimon_all_types全类型 可空/非空 多分区03546_paimon_all_supported_type.sql、03775_paimon_all_supported_type_storage.sqlpaimon_no_partition无分区键的简单表基础读取paimon_nullable_composites可空 ARRAY/MAP 复合类型04757_paimon_nullable_composite_types.sqlpaimon_timestamp_partitionTIMESTAMP 作为分区键04700_paimon_timestamp_partition.sql以paimon_nullable_composites为例04757_paimon_nullable_composite_types.sql 注释明确说明了这类数据集的用途可空的 Paimon ARRAY/MAP 列不能被包装成 Nullable否则会导致整表不可读对应 ClickHouse issue 113337。这说明生成器在设计 schema 时需刻意制造边界情况才能让消费端测试覆盖真实故障场景。七、注意事项与常见问题版本一致性pom 中paimon-flink-common为1.1-SNAPSHOT与paimon-common/paimon-core的${project.version}1.1.1不同源拉取失败时请检查仓库是否发布了对应用户/快照JDK 版本exec-maven-plugin的javagoal 在 JVM 中直接执行需确保JAVA_HOME指向 JDK 17README 实测为 17.0.12 LTSHadoop 运行时hadoop-common/hadoop-hdfs-client仅为runtime作用域若在 IDE 中直接运行主类请确认 IDE 的 classpath 配置包含了 runtime 依赖分区键选择将可空列选作分区键会生成__DEFAULT_PARTITION__目录仓库数据中可空分区列f_boolean/f_string/f_date的默认分区即如此命名这是 Paimon 的既定行为ClickHouse 侧通过partition.default-name选项对齐数据再生成README 的main中generateTestCase1(/tmp/warehouse)每次全量重建tests.cases2表如需生成仓库中paimon_all_types那样的数据集将rootPath指向目标 MinIO bucket 对应前缀并把数据上传到对象存储即可被paimonS3(s3_conn, filenamepaimon_all_types)读取。结语Paimon 格式目录的生成是 ClickHouse 数据湖测试的基础工程能力pom.xml负责锁定 Paimon/Flink/Hadoop 依赖矩阵DataGenerator负责构造覆盖全部类型、可空性与分区边界的真实数据而构建产物经mvn exec:java落盘后即成为paimonS3表函数与PaimonS3存储引擎的测试输入。理解这条生成 → 落盘 → 读取的完整链路无论是为 ClickHouse 扩展新的 Paimon 测试用例还是排查数据湖读取问题都能做到有据可依。赞分享数据库OLAP列式数据库大数据实时分析数据分析【免费下载链接】ClickHouseClickHouse® is a real-time analytics database management system项目地址https://gitcode.com/GitHub_Trending/cli/ClickHouse点击查看免费下载相关推荐使用 Paimon Java Client 生成 Paimon 无分区表目录并接入 ClickHouse 查询使用 Paimon Java Client 生成 Paimon 无分区表目录并接入 ClickHouse 查询 导读 Apache PaimonFlink 生数据库OLAP列式数据库大数据实时分析数据分析StarRocks Paimon Catalog 完整指南免数据导入直查 Apache Paimon 湖仓数据StarRocks Paimon Catalog 完整指南免数据导入直查 Apache Paimon 湖仓数据 本篇指南系统讲解如何在 StarRocks 中数据库OLAP数据仓库大数据湖仓一体数据分析ClickHouse 读取以 TIMESTAMP 为分区键的 Paimon 表分区目录重构原理与测试数据集生成指南ClickHouse 读取以 TIMESTAMP 为分区键的 Paimon 表分区目录重构原理与测试数据集生成指南 本文聚焦 ClickHouse 通过 pa数据库OLAP列式数据库大数据实时分析数据分析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考