开源湖仓一体平台一LakeSoul
LakeSoul 简介
LakeSoul 是由数元灵科技研发的云原生湖仓一体框架,具备高可扩展的元数据管理、ACID 事务、高效灵活的 upsert 操作、Schema 演进和批流一体化处理等特性。
主要特性
- 弹性架构:计算存储完全分离,不需要固定节点和磁盘,计算存储各自弹性扩容。并且针对云存储做了大量优化,在对象存储上实现了并发一致性、增量更新等功能;使用 LakeSoul 不需要维护固定的存储节点,云上对象存储的成本只有本地磁盘的 1/10,极大地降低了存储成本和运维成本。
- 高效可扩展的元数据管理:LakeSoul 使用 Postgres 数据库来管理文件元数据,可以高效的处理元数据的修改,并能够支持多并发写入,解决了 Hive 等元数据层的性能瓶颈,如长时间运行后元数据解析缓慢的痛点。元数据层的表结构经过精心设计,所有读写操作都能够使用主键索引,达到很高的 Ops。同时,元数据库在云上也能够很容易地进行扩容。
- ACID事务:通过元数据库事务机制实现了两阶段提交协议,保证了流批一体提交的事务性,用户不会看到不一致数据;支持多并发写入,自动冲突处理机制。
- 多级分区模式和高效灵活的upsert操作:LakeSoul 支持 range 和 hash 分区,通过灵活的 upsert 功能,支持行、列级别的增、删、改等更新操作,将 upsert 数据以 delta file 的形式保存,大幅提高了写数据效率和并发性,而优化过的 merge scan 提供了高效的 MergeOnRead 读取性能。
- 批流一体:LakeSoul 支持 streaming sink,可以同时处理流式数据摄入和历史数据批量回填、交互式查询等场景。
- Schema演进:支持新增、删除列,并在读取时自动兼容旧数据。
- CDC流、日志流自动同步:支持 MySQL 整库千表同步,自动建表和自动 Schema 变更;支持 Kafka 多 topic 合并同步、自动 Schema 解析、自动新 Topic 感知。
- 云对象存储IO优化:使用 Rust Arrow 实现原生 Parquet IO,并对对象存储访问做了专门优化。
适用场景
- 构建实时湖仓,并且新增数据需要高效实时大批量写入,同时需要行、列级别的并发增量更新的场景。
- 历史数据存储量很大,并且需要对大跨度时间范围做明细查询、修改,同时希望使用对象存储控制成本的场景。
- 查询请求不固定,资源消耗变化较大,希望计算资源能够独立弹性伸缩的场景。
- 需要多并发写,同时文件数量多,对元数据性能和并发有较高要求的场景。
- 针对主键进行数据更新,对写吞吐有较高有求的场景。
Spark + LakeSoul CDC 入湖
使用 Spark Streaming,消费 Kafka 数据并同步更新至 LakeSoul
1. 配置 LakeSoul 元数据库
自 2.1.0 起,LakeSoul Spark 和 Flink 的 jar 包通过 shade 方式打包了 Postgres Driver,Driver 的名字是 com.lakesoul.shaded.org.postgresql.Driver,而在 2.0.1 版本之前,Driver 还没有 shaded 打包,名字是 org.postgresql.Driver。
使用 LakeSoul 之前还需要初始化元数据表结构:
bash
PGPASSWORD=lakesoul_test psql -h localhost -p 5432 -U lakesoul_test -f script/meta_init.sql
LakeSoul 使用 lakesoul_home(大小写均可)环境变量或者 lakesoul_home JVM Property(只能全小写)来定位元数据库的配置文件,配置文件中主要包含 PostgreSQL DB 的连接信息。一个示例配置文件:
properties
lakesoul.pg.driver=com.lakesoul.shaded.org.postgresql.Driver
lakesoul.pg.url=jdbc:postgresql://localhost:5432/lakesoul_test?stringtype=unspecified
lakesoul.pg.username=lakesoul_test
lakesoul.pg.password=lakesoul_test
如果找不到上述环境变量或 JVM Property,则会分别查找 LAKESOUL_PG_DRIVER、LAKESOUL_PG_URL、LAKESOUL_PG_USERNAME、LAKESOUL_PG_PASSWORD 这几个环境变量作为配置的值。
2. 设置 Spark 工程作业
LakeSoul 目前支持 Spark 3.1.2 + Scala 2.12。
使用 spark-shell、pyspark 或者 spark-sql 交互式查询,需要添加 LakeSoul 的依赖和配置,有两种方法:
使用 --packages 传 Maven 仓库和包名
bash
spark-shell --packages com.dmetasoul:lakesoul-spark:2.1.1-spark-3.1.2
使用打包好的 LakeSoul 包
可以从 Releases 页面下载已经打包好的 LakeSoul Jar 包。下载 jar 并传给 spark-submit 命令:
bash
spark-submit --jars "lakesoul-spark-2.1.1-spark-3.1.2.jar"
将 Jar 包放在 Spark 环境中
将 Jar 包下载后,放在 $SPARK_HOME/jars 中。
Maven 依赖
xml
<dependency>
<groupId>com.dmetasoul</groupId>
<artifactId>lakesoul-spark</artifactId>
<version>2.1.1-spark-3.1.2</version>
</dependency>
3. Spark 作业设置 lakesoul_home 环境变量
bash
export lakesoul_home=/path/to/lakesoul.properties
- 对于 Hadoop Yarn 集群,增加命令行参数
--conf spark.yarn.appMasterEnv.lakesoul_home=lakesoul.properties --files /path/to/lakesoul.properties。 - 对于 K8s 集群,增加命令行参数
--conf spark.kubernetes.driverEnv.lakesoul_home=lakesoul.properties --files /path/to/lakesoul.properties到spark-submit命令。
4. 设置 Spark SQL Extension
LakeSoul 通过 Spark SQL Extension 机制来实现一些查询计划改写的扩展,需要为 Spark 作业增加以下配置:
bash
spark.sql.extensions=com.dmetasoul.lakesoul.sql.LakeSoulSparkSessionExtension
5. 设置 Spark 的 Catalog
LakeSoul 实现了 Spark 3 的 CatalogPlugin 接口,可以作为独立的 Catalog 插件让 Spark 加载。在 Spark 作业中增加如下配置:
bash
spark.sql.catalog.lakesoul=org.apache.spark.sql.lakesoul.catalog.LakeSoulCatalog
该配置增加了一个名为 lakesoul 的 Catalog。为了方便 SQL 中使用,也可以将该 Catalog 设置为默认的 Catalog:
bash
spark.sql.defaultCatalog=lakesoul
通过如上配置,默认会通过 LakeSoul Catalog 来查找所有 database 和表。如果需要同时访问 Hive 等外部 catalog,需要在表名前加上对应 catalog 名字。例如在 Spark 中启用 Hive 作为 Session Catalog,则访问 Hive 表时需要加上 spark_catalog 前缀。
bash
spark.sql.catalog.spark_catalog=org.apache.spark.sql.lakesoul.catalog.LakeSoulCatalog
从 2.1.0 起 LakeSoul 的 Catalog 更改为非 session 的实现。你仍然可以将 LakeSoul 设置为 Session Catalog,即设置名为 spark_catalog,但是这样就无法再访问到 Hive 表。
自 2.0 版本起,LakeSoul 支持将 Compaction 后的目录路径,挂载到指定的 Hive 表,指定和 LakeSoul 分区名一致和自定义分区名两种功能。该功能可以方便下游一些只能支持访问 Hive 的系统读取到 LakeSoul 的数据。更推荐的方式是通过 Kyuubi 来支持 Hive JDBC,这样可以直接使用 Hive JDBC 调用 Spark 引擎来访问 LakeSoul 表,包括 Merge on Read 读取。
scala
lakeSoulTable.compaction("date='2021-01-02'", "spark_catalog.default.hive_test_table", "date='20210102'")
6. 启动 Spark Shell
bash
./bin/spark-shell --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2 --conf spark.sql.extensions=com.dmetasoul.lakesoul.sql.LakeSoulSparkSessionExtension --conf spark.sql.catalog.spark_catalog=org.apache.spark.sql.lakesoul.catalog.LakeSoulCatalog
7. 创建 LakeSoul 表
我们创建一个 LakeSoul 表 MysqlCdcTest,这个表会准实时流批一体地同步 MySQL 的数据。这个表同样使用 id 列作为主键,用 "op" 列表示 CDC 的更新。并且我们需要通过 lakesoul_cdc_change_column 这个表属性,指定 LakeSoul 表中,表示 CDC 状态更新的列名,这个示例中该列的名字为 "op"。
scala
import com.dmetasoul.lakesoul.tables.LakeSoulTable
val path = "/opt/spark/cdctest"
val data = Seq((1L, 1L, "hello world", "insert")).toDF("id", "rangeid", "value", "op")
LakeSoulTable.createTable(data, path)
.shortTableName("cdc")
.hashPartitions("id")
.hashBucketNum(2)
.rangePartitions("rangeid")
.tableProperty("lakesoul_cdc_change_column" -> "op")
.create()
8. 启动 Streaming 写入 LakeSoul
读取 Kafka,转换 Debezium 读取出的 JSON 格式,使用 lakesoul_upsert 更新 LakeSoul 表:
scala
import com.dmetasoul.lakesoul.tables.LakeSoulTable
val path = "/opt/spark/cdctest"
val lakeSoulTable = LakeSoulTable.forPath(path)
var strList = List.empty[String]
// js1 是示例数据,我们这里也用于生成 schema,在下文 from_json 函数中转换数据使用,before 和 after 中内容对应于 mysql 表字段
val js1 = """
{
"before": {
"id": 2,
"rangeid": 2,
"value": "sms"
},
"after": {
"id": 2,
"rangeid": 2,
"value": "sms"
},
"source": {
"version": "1.8.0.Final",
"connector": "mysql",
"name": "cdcserver",
"ts_ms": 1644461444000,
"snapshot": "false",
"db": "cdc",
"sequence": null,
"table": "sms",
"server_id": 529210004,
"gtid": "de525a81-57f6-11ec-9b60-fa163e692542:1621099",
"file": "binlog.000033",
"pos": 54831329,
"row": 0,
"thread": null,
"query": null
},
"op": "c",
"ts_ms": 1644461444777,
"transaction": null
}
""".stripMargin
strList = strList :+ js1
val rddData = spark.sparkContext.parallelize(strList)
val resultDF = spark.read.json(rddData)
val sche = resultDF.schema
import org.apache.spark.sql.{DataFrame, SaveMode, SparkSession}
// 对接 Kafka 需要指定 kafka.bootstrap.server ip 地址和 debezium 输出到 Kafka 的 topic
val kfdf = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "kafkahost:9092")
.option("subscribe", "cdcserver.cdc.test")
.option("startingOffsets", "latest")
.load()
// 解析 debezium 中产生的 JSON,对于 MySQL insert、update、delete 操作会自动生成一列,列名 op
val kfdfdata = kfdf
.selectExpr("CAST(value AS STRING) as value")
.withColumn("payload", from_json($"value", sche))
.filter("value is not null")
.drop("value")
.select("payload.after", "payload.before", "payload.op")
.withColumn(
"op",
when($"op" === "c", "insert")
.when($"op" === "u", "update")
.when($"op" === "d", "delete")
.otherwise("unknown")
)
.withColumn(
"data",
when($"op" === "insert" || $"op" === "update", $"after")
.when($"op" === "delete", $"before")
)
.drop($"after")
.drop($"before")
.select("data.*", "op")
// 使用 lakesoul upsert 更新表中数据并在屏幕上输出解析后的数据
kfdfdata.writeStream
.foreachBatch { (batchDF: DataFrame, _: Long) =>
{
lakeSoulTable.upsert(batchDF)
batchDF.show
}
}
.start()
.awaitTermination()
Flink CDC 同步到 LakeSoul
LakeSoul 自 2.1.0 版本起,实现了 Flink CDC Sink,能够支持 Table API 及 SQL(单表),以及 Stream API(整库多表)。目前支持的上游数据源为 MySQL(5.6-8.0)。
1. 下载 LakeSoul Flink Jar
下载 lakesoul-flink-2.1.1-flink-1.14.jar。
2. 增加 LakeSoul 元数据库配置
在 $FLINK_HOME/conf/flink-conf.yaml 中增加如下配置:
yaml
containerized.master.env.LAKESOUL_PG_DRIVER: com.lakesoul.shaded.org.postgresql.Driver
containerized.master.env.LAKESOUL_PG_USERNAME: root
containerized.master.env.LAKESOUL_PG_PASSWORD: root
containerized.master.env.LAKESOUL_PG_URL: jdbc:postgresql://192.168.24.180:5432/test_lakesoul_meta?stringtype=unspecified
containerized.taskmanager.env.LAKESOUL_PG_DRIVER: com.lakesoul.shaded.org.postgresql.Driver
containerized.taskmanager.env.LAKESOUL_PG_USERNAME: root
containerized.taskmanager.env.LAKESOUL_PG_PASSWORD: root
containerized.taskmanager.env.LAKESOUL_PG_URL: jdbc:postgresql://192.168.24.180:5432/test_lakesoul_meta?stringtype=unspecified
注意这里 master 和 taskmanager 的环境变量都需要设置。
如果使用 Session 模式来启动作业,即将作业以 client 方式提交到 Flink Standalone Cluster,则 flink run 作为 client,是不会读取上面配置,因此需要再单独配置环境变量,即:
bash
export LAKESOUL_PG_DRIVER=com.lakesoul.shaded.org.postgresql.Driver
export LAKESOUL_PG_URL=jdbc:postgresql://localhost:5432/test_lakesoul_meta?stringtype=unspecified
export LAKESOUL_PG_USERNAME=root
export LAKESOUL_PG_PASSWORD=root
3. 启动同步作业
bash
bin/flink run -c org.apache.flink.lakesoul.entry.MysqlCdc
lakesoul-flink-2.1.1-flink-1.14.jar
--source_db.host localhost
--source_db.port 3306
--source_db.db_name default
--source_db.user root
--source_db.password root
--source.parallelism 4
--sink.parallelism 4
--server_time_zone=Asia/Shanghai
--warehouse_path s3://bucket/lakesoul/flink/data
--flink.checkpoint s3://bucket/lakesoul/flink/checkpoints
--flink.savepoint s3://bucket/lakesoul/flink/savepoints
4. DataStream 代码
java
package org.apache.flink.lakesoul.entry;
import com.dmetasoul.lakesoul.meta.external.mysql.MysqlDBManager;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.connectors.mysql.source.MySqlSourceBuilder;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.lakesoul.sink.LakeSoulMultiTableSinkStreamBuilder;
import org.apache.flink.lakesoul.tool.LakeSoulSinkOptions;
import org.apache.flink.lakesoul.types.BinarySourceRecord;
import org.apache.flink.streaming.api.CheckpointingMode;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSink;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.CheckpointConfig;
import org.apache.flink.streaming.api.environment.ExecutionCheckpointingOptions;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.util.HashSet;
import java.util.List;
import static org.apache.flink.lakesoul.tool.JobOptions.*;
import static org.apache.flink.lakesoul.tool.LakeSoulDDLSinkOptions.*;
public class MysqlCdc2LakeSoul {
public static void main(String[] args) throws Exception {
ParameterTool parameter = ParameterTool.fromArgs(args);
String dbName = parameter.get(SOURCE_DB_DB_NAME.key());
String userName = parameter.get(SOURCE_DB_USER.key());
String passWord = parameter.get(SOURCE_DB_PASSWORD.key());
String host = parameter.get(SOURCE_DB_HOST.key());
int port = parameter.getInt(SOURCE_DB_PORT.key(), MysqlDBManager.DEFAULT_MYSQL_PORT);
String databasePrefixPath = parameter.get(WAREHOUSE_PATH.key());
String serverTimezone = parameter.get(SERVER_TIME_ZONE.key(), SERVER_TIME_ZONE.defaultValue());
int sourceParallelism = parameter.getInt(SOURCE_PARALLELISM.key());
int bucketParallelism = parameter.getInt(BUCKET_PARALLELISM.key());
int checkpointInterval = parameter.getInt(JOB_CHECKPOINT_INTERVAL.key(), JOB_CHECKPOINT_INTERVAL.defaultValue()); // mill second
MysqlDBManager mysqlDBManager = new MysqlDBManager(dbName, userName, passWord, host, Integer.toString(port), new HashSet<>(), databasePrefixPath, bucketParallelism, true);
mysqlDBManager.importOrSyncLakeSoulNamespace(dbName); // syncing mysql tables to lakesoul
List<String> tableList = mysqlDBManager.listTables();
if (tableList.isEmpty()) {
throw new IllegalStateException("Failed to discover captured tables");
}
tableList.forEach(mysqlDBManager::importOrSyncLakeSoulTable);
Configuration conf = new Configuration();
// parameters for mutil tables ddl sink
conf.set(SOURCE_DB_DB_NAME, dbName);
conf.set(SOURCE_DB_USER, userName);
conf.set(SOURCE_DB_PASSWORD, passWord);
conf.set(SOURCE_DB_HOST, host);
conf.set(SOURCE_DB_PORT, port);
conf.set(WAREHOUSE_PATH, databasePrefixPath);
conf.set(SERVER_TIME_ZONE, serverTimezone);
// parameters for mutil tables dml sink
conf.set(LakeSoulSinkOptions.USE_CDC, true);
conf.set(LakeSoulSinkOptions.WAREHOUSE_PATH, databasePrefixPath);
conf.set(LakeSoulSinkOptions.SOURCE_PARALLELISM, sourceParallelism);
conf.set(LakeSoulSinkOptions.BUCKET_PARALLELISM, bucketParallelism);
StreamExecutionEnvironment env;
env = StreamExecutionEnvironment.getExecutionEnvironment();
ParameterTool pt = ParameterTool.fromMap(conf.toMap());
env.getConfig().setGlobalJobParameters(pt);
env.enableCheckpointing(checkpointInterval);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(4023);
CheckpointingMode checkpointingMode = CheckpointingMode.EXACTLY_ONCE;
if (parameter.get(JOB_CHECKPOINT_MODE.key(), JOB_CHECKPOINT_MODE.defaultValue()).equals("AT_LEAST_ONCE")) {
checkpointingMode = CheckpointingMode.AT_LEAST_ONCE;
}
env.getCheckpointConfig().setCheckpointingMode(checkpointingMode);
env.getCheckpointConfig().setExternalizedCheckpointCleanup(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
env.getCheckpointConfig().setCheckpointStorage(parameter.get(FLINK_CHECKPOINT.key()));
conf.set(ExecutionCheckpointingOptions.ENABLE_CHECKPOINTS_AFTER_TASKS_FINISH, true);
MySqlSourceBuilder<BinarySourceRecord> sourceBuilder = MySqlSource.<BinarySourceRecord>builder()
.hostname(host)
.port(port)
.databaseList(dbName) // set captured database
.tableList(dbName + ".*") // set captured table
.serverTimeZone(serverTimezone) // default -- Asia/Shanghai
.username(userName)
.password(passWord);
LakeSoulMultiTableSinkStreamBuilder.Context context = new LakeSoulMultiTableSinkStreamBuilder.Context();
context.env = env;
context.sourceBuilder = sourceBuilder;
context.conf = conf;
LakeSoulMultiTableSinkStreamBuilder builder = new LakeSoulMultiTableSinkStreamBuilder(context);
DataStreamSource<BinarySourceRecord> source = builder.buildMultiTableSource();
Tuple2<DataStream<BinarySourceRecord>, DataStream<BinarySourceRecord>> streams = builder.buildCDCAndDDLStreamsFromSource(source);
DataStream<BinarySourceRecord> stream = builder.buildHashPartitionedCDCStream(streams.f0);
DataStreamSink<BinarySourceRecord> dmlSink = builder.buildLakeSoulDMLSink(stream);
DataStreamSink<BinarySourceRecord> ddlSink = builder.buildLakeSoulDDLSink(streams.f1);
dmlSink.addSink(new LakeSoulMultiTablesSink());
ddlSink.addSink(new LakeSoulDDLSink());
env.execute("LakeSoul CDC Sink From MySQL Database " + dbName);
}
}
5. LakeSoul Flink CDC Sink 严格一次语义保证
LakeSoul Flink CDC Sink 在作业运行过程中会自动保存相关状态,在 Flink 作业发生 Failover 时能够将状态恢复并重新写入,因此数据不会丢失。
LakeSoul 写入时,在两个部分保证写入的幂等性:
- Stage 文件 Commit 时:与 Flink File Sink 一致,通过文件系统 rename 操作的原子性,来保证 staging 文件写入到最终的路径。因为 rename 是原子的,Failover 之后不会发生重复写入或缺失的情况。
- LakeSoul 元数据提交时:会首先记录文件路径,在更新 snapshot 时会通过事务标记该文件已提交。Failover 后,通过判断一个文件是否已经提交,可以保证提交的幂等性。
综上,LakeSoul Flink CDC Sink 通过状态恢复保证数据不丢失,通过提交幂等性保证数据不重复,实现了严格一次(Exactly Once)语义保证。
数据更新和合并(Upsert/Merge)
LakeSoul 可以支持对已经入湖的数据做 部分字段更新功能,而不必将整张数据表全部覆盖重写,避免这种繁重且浪费资源的操作。
1. 部分字段更新
举个例子一张表数据信息如下,id 为主键(即 hashPartitions),目前需要根据主键字段,对 phone_number 做字段修改处理。

可以使用 upsert 来实现对任意行中任意一个字段的更新。upsert 需要包含主键 (id) 和需要修改的 address 信息,再次读取整张表数据 address 便可展示为修改后的字段信息。
scala
import org.apache.spark.sql._
val spark = SparkSession.builder.master("local")
.config("spark.sql.extensions", "com.dmetasoul.lakesoul.sql.LakeSoulSparkSessionExtension")
.getOrCreate()
import spark.implicits._
val df = Seq(
("1", "Jake", "13700001111", "address_1", "job_1", "company_1"),
("2", "Make", "13511110000", "address_2", "job_2", "company_2")
).toDF("id", "name", "phone_number", "address", "job", "company")
val tablePath = "s3a://bucket-name/table/path/is/also/table/name"
df.write
.mode("append")
.format("lakesoul")
.option("hashPartitions", "id")
.option("hashBucketNum", "2")
.save(tablePath)
val lakeSoulTable = LakeSoulTable.forPath(tablePath)
val extraDF = Seq(("1", "address_1_1")).toDF("id", "address")
lakeSoulTable.upsert(extraDF)
lakeSoulTable.toDF.show()

2. 自定义 Merge 合并功能
LakeSoul 默认 merge 规则,即 数据更新后取最后一条记录作为该字段数据(org.apache.spark.sql.execution.datasources.v2.merge.parquet.batch.merge_operator.DefaultMergeOp)。在此基础上,LakeSoul 内置扩展了几种数据 merge 逻辑,对 Int/Long 字段做加和 merge(MergeOpInt/MergeOpLong)、对非空字段更新(MergeNonNullOp)、以 "," 拼接字符串 merge 方式。
下面以对非空字段更新(MergeNonNullOp)为例,借用上面表格数据样例。数据写入时同样以 upsert 方式进行更新写入,然后在数据读取时需要注册 merger 逻辑,然后进行读取即可。
scala
import org.apache.spark.sql.execution.datasources.v2.merge.parquet.batch.merge_operator.MergeNonNullOp
import org.apache.spark.sql.functions.expr
import org.apache.spark.sql._
val spark = SparkSession.builder.master("local")
.config("spark.sql.extensions", "com.dmetasoul.lakesoul.sql.LakeSoulSparkSessionExtension")
.getOrCreate()
import spark.implicits._
val df = Seq(
("1", "Jake", "13700001111", "address_1", "job_1", "company_1"),
("2", "Make", "13511110000", "address_2", "job_2", "company_2")
).toDF("id", "name", "phone_number", "address", "job", "company")
val tablePath = "s3a://bucket-name/table/path/is/also/table/name"
df.write
.mode("append")
.format("lakesoul")
.option("hashPartitions", "id")
.option("hashBucketNum", "2")
.save(tablePath)
val lakeSoulTable = LakeSoulTable.forPath(tablePath)
val extraDF = Seq(
("1", "null", "13100001111", "address_1_1", "job_1_1", "company_1_1"),
("2", "null", "13111110000", "address_2_2", "job_2_2", "company_2_2")
).toDF("id", "name", "phone_number", "address", "job", "company")
new MergeNonNullOp().register(spark, "NotNullOp")
lakeSoulTable.toDF.show()
lakeSoulTable.upsert(extraDF)
lakeSoulTable.toDF.withColumn("name", expr("NotNullOp(name)")).show()

用户也可以通过自定义 MergeOperator(实现 trait org.apache.spark.sql.execution.datasources.v2.merge.parquet.batch.merge_operator.MergeOperator)来自定义 Merge 时的逻辑,能够灵活地实现数据高效入湖。
多流合并构建宽表
为构建宽表,传统数仓的 ETL 在做多表关联时,需要根据主外键多次 join,然后构建一个大宽表。当数据量较多或需要多次 join 时,会有效率低下,内存消耗大,容易 OOM 等问题,且 Shuffle 过程占据大部分数据交换时间,效率也很低下。LakeSoul 支持对数据进行 Upsert,并支持自定义 MergeOperator 功能,可以避免上述存在的问题,不必 Join 即可得到合并结果。下面针对这一场景具体举例进行说明。
假设有以下几个流的数据,A、B、C 和 D,各个流数据内容如下:
A:

B:

C:

D:

最后需要形成一张大宽表,将四张表进行合并展示,如下:

传统意义上进行上述操作,需要将四张表根据主键(IP)进行三次 join,写法如下:
sql
SELECT
A.IP AS IP,
A.sy AS sy,
A.us AS us,
B.free AS free,
B.cache AS cache,
C.level AS level,
C.des AS des,
D.qps AS qps,
D.tps AS tps
FROM A
JOIN B ON A.IP = B.IP
JOIN C ON C.IP = A.IP
JOIN D ON D.IP = A.IP
LakeSoul 支持多流合并,多个流可以有不同的 Schema(需要有相同主键)。LakeSoul 可以做到自动扩展 Schema,若新写入的数据字段在原表中未存在,则会自动扩展表 schema,不存在的字段默认为 null 处理。通过使用 LakeSoul 多流合并功能,结合 LakeSoul 独特的 MergeOperator 功能,通过 upsert 将数据写入 LakeSoul 后,不需要 join,即可读取到拼接好的宽表。上述过程代码实现如下:
scala
import org.apache.spark.sql._
val spark = SparkSession.builder.master("local")
.config("spark.sql.extensions", "com.dmetasoul.lakesoul.sql.LakeSoulSparkSessionExtension")
.config("spark.dmetasoul.lakesoul.schema.autoMerge.enabled", "true")
.getOrCreate()
import spark.implicits._
val df1 = Seq(("1.1.1.1", 30, 40)).toDF("IP", "sy", "us")
val df2 = Seq(("1.1.1.1", 1677, 455)).toDF("IP", "free", "cache")
val df3 = Seq(("1.1.1.2", "error", "killed")).toDF("IP", "level", "des")
val df4 = Seq(("1.1.1.1", 30, 40)).toDF("IP", "qps", "tps")
val tablePath = "s3a://bucket-name/table/path/is/also/table/name"
df1.write
.mode("append")
.format("lakesoul")
.option("hashPartitions", "IP")
.option("hashBucketNum", "2")
.save(tablePath)
val lakeSoulTable = LakeSoulTable.forPath(tablePath)
lakeSoulTable.upsert(df2)
lakeSoulTable.upsert(df3)
lakeSoulTable.upsert(df4)
lakeSoulTable.toDF.show()

Kafka 多 Topic 入 LakeSoul
通过 LakeSoul Kafka Stream 将 Kafka 中的数据同步到 LakeSoul 非常方便。LakeSoul Kafka Stream 可以支持自动创建表,自动识别新 topic,exactly-once 语义、自动为表添加分区等功能。LakeSoul Kafka Stream 主要使用 Spark Structured Streaming 来实现数据同步功能。
使用 LakeSoul Kafka Stream 需要以下条件之一:
- topic 中的数据为 JSON 格式;
- Kafka 集群带有 Schema Registry 服务。
1. 准备环境
你可以编译 LakeSoul 项目以获取 LakeSoul Kafka Stream jar,或者可以通过 LakeSoul Kafka Stream 来获取 LakeSoul Kafka Stream 以及其他任务运行依赖的 jar 包。
下载后解压 tar 包,然后将 jar 包放入 $SPARK_HOME/jars 目录下,或者在提交任务时添加依赖的 jar,比如通过 --jars。
2. 启动 LakeSoul Kafka Stream 任务
- 任务启动时通过
lakesoul_home环境变量添加元数据库信息。这部分请参考 搭建本地测试环境。 - 提交任务。你需要按顺序填写一些参数,以确保任务能够准确运行。参数描述如下:

3. 任务流程示例
- 假设 Kafka 集群已经存在。在这里,通过 Docker Compose 运行 Kafka 集群。然后创建一个名为 "test" 的主题并向其中写入一些数据。Kafka bootstrap.servers:
localhost:9092
bash
# 创建 topic 'test'
bin/kafka-topics.sh --create --topic test --bootstrap-server localhost:9092 --replication-factor 1 --partitions 4
# 查看 topic 列表
bin/kafka-topics.sh --list --bootstrap-server localhost:9092
test
# 向
end
