Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 4 additions & 7 deletions bubble-dependencies/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@
<forest.version>1.8.0</forest.version>
<java-diff-utils.version>4.15</java-diff-utils.version>
<liteflow.version>2.15.2</liteflow.version>
<duckdb.version>1.4.2.0</duckdb.version>
</properties>

<dependencyManagement>
Expand Down Expand Up @@ -113,19 +114,19 @@

<dependency>
<groupId>cn.fxbin.bubble</groupId>
<artifactId>bubble-starter-lock</artifactId>
<artifactId>bubble-starter-data-duckdb</artifactId>
<version>2.0.0.BUILD-SNAPSHOT</version>
</dependency>

<dependency>
<groupId>cn.fxbin.bubble</groupId>
<artifactId>bubble-starter-excel</artifactId>
<artifactId>bubble-starter-lock</artifactId>
<version>2.0.0.BUILD-SNAPSHOT</version>
</dependency>

<dependency>
<groupId>cn.fxbin.bubble</groupId>
<artifactId>bubble-starter-flow</artifactId>
<artifactId>bubble-starter-excel</artifactId>
<version>2.0.0.BUILD-SNAPSHOT</version>
</dependency>

Expand Down Expand Up @@ -525,9 +526,6 @@
<artifactId>forest-spring-boot3-starter</artifactId>
<version>${forest.version}</version>
</dependency>
<<<<<<< HEAD

=======

<!-- duckdb -->
<dependency>
Expand All @@ -536,7 +534,6 @@
<version>${duckdb.version}</version>
</dependency>

>>>>>>> df135a0 (feat(bubble-dependencies): 引入liteflow规则引擎依赖)
<!-- java-diff-utils -->
<dependency>
<groupId>io.github.java-diff-utils</groupId>
Expand Down
45 changes: 45 additions & 0 deletions bubble-starters/bubble-starter-data-duckdb/pom.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>cn.fxbin.bubble</groupId>
<artifactId>bubble-starters</artifactId>
<version>2.0.0.BUILD-SNAPSHOT</version>
</parent>

<artifactId>bubble-starter-data-duckdb</artifactId>
<name>Bubble Starter Data DuckDB</name>

<dependencies>
<dependency>
<groupId>cn.fxbin.bubble</groupId>
<artifactId>bubble-starter</artifactId>
</dependency>

<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>

<dependency>
<groupId>org.duckdb</groupId>
<artifactId>duckdb_jdbc</artifactId>
</dependency>

<dependency>
<groupId>net.dreamlu</groupId>
<artifactId>mica-auto</artifactId>
<scope>provided</scope>
</dependency>

<!-- Testing dependencies -->
<dependency>
<groupId>cn.fxbin.bubble</groupId>
<artifactId>bubble-starter-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>

</project>
Original file line number Diff line number Diff line change
@@ -0,0 +1,200 @@
package cn.fxbin.bubble.data.duckdb;

import cn.fxbin.bubble.data.duckdb.core.DuckDbIngester;
import cn.fxbin.bubble.data.duckdb.core.DuckDbManager;
import cn.fxbin.bubble.data.duckdb.core.DuckDbTemplate;
import lombok.extern.slf4j.Slf4j;
import org.springframework.util.Assert;

import java.util.Iterator;
import java.util.List;
import java.util.Map;

/**
* DuckDB 统一客户端入口
*
* <p>
* 提供对默认 DuckDB 实例的直接访问,以及对动态 DuckDB 文件的管理。
* 采用了外观模式(Facade Pattern),将 DuckDbTemplate、DuckDbIngester 和 DuckDbManager 的功能统一暴露。
* </p>
*
* @author fxbin
* @version v1.0
* @since 2025/12/08 16:30
*/
@Slf4j
public class DuckDbOperations {

private final DuckDbTemplate defaultTemplate;
private final DuckDbManager manager;
private final DuckDbIngester defaultIngester;

public DuckDbOperations(DuckDbTemplate defaultTemplate, DuckDbManager manager, DuckDbIngester defaultIngester) {
Assert.notNull(defaultTemplate, "Default DuckDbTemplate must not be null");
Assert.notNull(manager, "DuckDbManager must not be null");
Assert.notNull(defaultIngester, "Default DuckDbIngester must not be null");
this.defaultTemplate = defaultTemplate;
this.manager = manager;
this.defaultIngester = defaultIngester;
}

// =================================================================================================================
// 基础查询操作 (Delegate to defaultTemplate)
// =================================================================================================================

/**
* 获取默认的 DuckDbTemplate。
*
* @return 默认模板实例。
*/
public DuckDbTemplate template() {
return defaultTemplate;
}

/**
* 在默认数据库上执行 SQL 查询。
*
* @param sql SQL 查询语句。
* @return 结果列表(Map)。
*/
public List<Map<String, Object>> query(String sql) {
return defaultTemplate.queryForList(sql);
}

/**
* 在默认数据库上执行 SQL 查询,并返回指定类型的对象列表。
*
* @param sql SQL 查询语句。
* @param type 结果类型。
* @param <T> 泛型类型。
* @return 对象列表。
*/
public <T> List<T> query(String sql, Class<T> type) {
return defaultTemplate.queryForList(sql, type);
}

/**
* 在默认数据库上执行 SQL 查询,并返回单个对象。
*
* @param sql SQL 查询语句。
* @param type 结果类型。
* @param <T> 泛型类型。
* @return 单个对象,如果未找到则为 null。
*/
public <T> T queryForObject(String sql, Class<T> type) {
try {
return defaultTemplate.queryForObject(sql, type);
} catch (org.springframework.dao.EmptyResultDataAccessException e) {
return null;
}
}

/**
* 获取表行数统计。
*
* @param tableName 表名。
* @return 行数。
*/
public Long count(String tableName) {
DuckDbTemplate.validateTableName(tableName);
String sql = "SELECT count(*) FROM " + tableName;
return defaultTemplate.queryForObject(sql, Long.class);
}

/**
* 在默认数据库上执行 SQL 语句(DDL/DML)。
*
* @param sql SQL 语句。
*/
public void execute(String sql) {
defaultTemplate.execute(sql);
}

// =================================================================================================================
// 数据导入导出 (Ingestion & Parquet)
// =================================================================================================================

/**
* 向默认数据库追加数据(使用 Appender,高性能)。
*
* @param tableName 表名。
* @param dataIterator 数据迭代器。
*/
public void ingest(String tableName, Iterator<Object[]> dataIterator) {
defaultIngester.ingest(tableName, dataIterator);
}

/**
* 向默认数据库追加数据(列表方式)。
*
* @param tableName 表名。
* @param rows 数据行。
*/
public void append(String tableName, List<Object[]> rows) {
defaultTemplate.append(tableName, rows);
}

/**
* 将 Parquet 文件导入到表中。
*
* @param tableName 目标表名(将被创建)。
* @param parquetPath Parquet 文件的路径。
*/
public void importParquet(String tableName, String parquetPath) {
defaultTemplate.importParquet(tableName, parquetPath);
}

/**
* 将表(或查询结果)导出为 Parquet 文件。
*
* @param tableNameOrQuery 表名或 SELECT 查询。
* @param outputPath Parquet 文件的输出路径。
*/
public void exportParquet(String tableNameOrQuery, String outputPath) {
defaultTemplate.exportParquet(tableNameOrQuery, outputPath);
}

// =================================================================================================================
// 动态实例操作 (Delegate to manager)
// =================================================================================================================

/**
* 连接到指定的 DuckDB 文件。
*
* @param filePath 数据库文件路径。
* @return 该文件的 DuckDbTemplate 实例。
*/
public DuckDbTemplate connect(String filePath) {
return manager.getTemplate(filePath);
}

/**
* 以指定模式连接到指定的 DuckDB 文件。
*
* @param filePath 数据库文件路径。
* @param readOnly 是否只读。
* @return 该文件的 DuckDbTemplate 实例。
*/
public DuckDbTemplate connect(String filePath, boolean readOnly) {
return manager.getTemplate(filePath, readOnly);
}

/**
* 关闭指定文件的连接。
*
* @param filePath 数据库文件路径。
*/
public void close(String filePath) {
manager.close(filePath, false);
manager.close(filePath, true); // 尝试关闭只读连接
}

/**
* 获取底层的 DuckDbManager。
*
* @return 管理器实例。
*/
public DuckDbManager manager() {
return manager;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
package cn.fxbin.bubble.data.duckdb.autoconfigure;

import cn.fxbin.bubble.data.duckdb.DuckDbOperations;
import cn.fxbin.bubble.data.duckdb.core.DuckDbConnectionFactory;
import cn.fxbin.bubble.data.duckdb.core.DuckDbIngester;
import cn.fxbin.bubble.data.duckdb.core.DuckDbManager;
import cn.fxbin.bubble.data.duckdb.core.DuckDbTemplate;
import com.zaxxer.hikari.HikariDataSource;
import lombok.extern.slf4j.Slf4j;
import org.duckdb.DuckDBDriver;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.AutoConfiguration;
import org.springframework.boot.autoconfigure.AutoConfigureAfter;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;

import javax.sql.DataSource;

/**
* DuckDB 自动配置
*
* @author fxbin
* @version v1.0
* @since 2025/12/08 11:42
*/
@Slf4j
@AutoConfiguration
@ConditionalOnClass({DuckDBDriver.class, DataSource.class})
@EnableConfigurationProperties(DuckDbProperties.class)
@ConditionalOnProperty(prefix = "dm.data.duckdb", name = "enabled", havingValue = "true", matchIfMissing = true)
@AutoConfigureAfter(DataSourceAutoConfiguration.class)
public class DuckDbAutoConfiguration {

@Bean
@ConditionalOnMissingBean
public DuckDbConnectionFactory duckDbConnectionFactory(DuckDbProperties properties) {
return new DuckDbConnectionFactory(properties);
}

/**
* 定义 DuckDB 专用数据源
* <p>
* 注意:这里不使用 @Primary,也不使用 @ConditionalOnMissingBean(DataSource.class),
* 以确保它能作为“第二数据源”与主业务库(如 MySQL)共存。
* </p>
*/
@Bean("duckDbDataSource")
public DataSource duckDbDataSource(DuckDbConnectionFactory connectionFactory, DuckDbProperties properties) {
if (properties.getMode() == DuckDbProperties.Mode.FILE && !properties.isReadOnly()) {
log.warn("DuckDB 配置为 FILE 模式,具有 READ_WRITE 访问权限。" +
"确保没有其他进程正在访问 '{}'。DuckDB 强制执行单一写入者策略。", properties.getFilePath());
}

HikariDataSource dataSource = connectionFactory.createDefaultDataSource();
log.info("DuckDB 数据源已初始化: {}", dataSource.getJdbcUrl());
return dataSource;
}

@Bean
@ConditionalOnMissingBean
public DuckDbTemplate duckDbTemplate(@Qualifier("duckDbDataSource") DataSource dataSource) {
return new DuckDbTemplate(dataSource);
}

@Bean
@ConditionalOnMissingBean
public DuckDbManager duckDbManager(DuckDbProperties properties) {
return new DuckDbManager(properties);
}

@Bean
@ConditionalOnMissingBean
public DuckDbIngester duckDbIngester(@Qualifier("duckDbDataSource") DataSource dataSource) {
return new DuckDbIngester(dataSource);
}

@Bean
@ConditionalOnMissingBean
public DuckDbOperations duckDbOperations(DuckDbTemplate duckDbTemplate, DuckDbManager duckDbManager, DuckDbIngester duckDbIngester) {
return new DuckDbOperations(duckDbTemplate, duckDbManager, duckDbIngester);
}

}
Loading