改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

开源湖仓一体平台一LakeSoul

ckckck2025年1月10日4 浏览

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_DRIVERLAKESOUL_PG_URLLAKESOUL_PG_USERNAMELAKESOUL_PG_PASSWORD 这几个环境变量作为配置的值。

2. 设置 Spark 工程作业

LakeSoul 目前支持 Spark 3.1.2 + Scala 2.12。

使用 spark-shellpyspark 或者 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.propertiesspark-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)。

下载 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);
    }
}

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 需要以下条件之一:

  1. topic 中的数据为 JSON 格式;
  2. 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 任务

  1. 任务启动时通过 lakesoul_home 环境变量添加元数据库信息。这部分请参考 搭建本地测试环境
  2. 提交任务。你需要按顺序填写一些参数,以确保任务能够准确运行。参数描述如下:

3. 任务流程示例

  1. 假设 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