Featured image of post DuckDB 纯 Java 表函数:零 Native 成本接入任意数据源

DuckDB 纯 Java 表函数:零 Native 成本接入任意数据源

DuckDB Java 客户端新增纯 Java 表函数能力,无需 C++ 扩展即可将 MongoDB、JDBC 等任何 Java 数据源暴露为 SQL 表,实现异构数据源的单节点联邦查询。

DuckDB 纯 Java 表函数架构

引言:数据孤岛时代的查询困境

在大型企业环境中,数据分散在关系数据库、文档存储、消息队列、数据湖和云数仓等数十个系统中。许多系统只能通过厂商提供的 Java SDK 访问——Oracle 的 JDBC、MongoDB 的 Java Driver、甚至内部的 SOAP 接口。要将这些数据关联分析,传统方案需要部署 Trino 等分布式查询引擎,或者先将数据导出到 Parquet 再加载到分析工具。

DuckDB 最近发布了一个革命性功能:纯 Java 表函数(Pure Java Table Functions)。这意味着你可以在 Java 应用中直接注册一个自定义表函数,将任何 Java 可访问的数据源暴露为 SQL 表,然后在 DuckDB 中与其他数据源进行联邦查询——无需编写一行 C++ 代码,无需构建 DuckDB 扩展。

本文基于 DuckDB 官方博客,通过 MongoDB 示例完整讲解这个功能。


背景:为什么需要纯 Java 表函数?

传统方案的问题

在没有纯 Java 表函数之前,要在 DuckDB 中访问一个自定义数据源,你有以下选择:

方案复杂度风险适用场景
导出为 Parquet/CSV数据过时、额外存储一次性分析
编写 C++ 扩展Segfault 风险、构建复杂高性能需求
部署 Trino/Presto基础设施复杂大规模分布式查询
用 JDBC 读取后加载内存压力大小数据集

纯 Java 表函数的优势

纯 Java 表函数解决了上述所有问题:

  1. 零 Native 依赖:项目是纯 Maven 依赖,无需 C++ 工具链
  2. 复用已有客户端:直接使用 vendor 的 Java Driver,不重新实现协议
  3. 流式读取:数据通过游标逐批流入 DuckDB,不占用全部内存
  4. 谓词下推:过滤条件以源方言传给驱动,远程执行索引查询
  5. 联邦查询:远程数据与本地 CSV/Parquet 在同一 SQL 中 JOIN

核心概念:DuckDB 表函数生命周期

DuckDB 的表函数遵循一个三阶段生命周期:

┌─────────────────────────────────────────────────────────────┐
│                    表函数生命周期                             │
├──────────────┬──────────────┬──────────────┬────────────────┤
│   BIND       │     INIT     │    APPLY     │   清理          │
│  (准备阶段)   │  (初始化阶段)  │  (执行阶段)   │                │
├──────────────┼──────────────┼──────────────┼────────────────┤
│ • 读取参数    │ • 打开游标    │ • 填充数据块  │ • 关闭连接      │
│ • 声明输出Schema│ • 建立连接   │ • 返回行数   │ • 释放资源      │
│ • 创建绑定对象 │              │   (≤2048行)  │                │
└──────────────┴──────────────┴──────────────┴────────────────┘
  • bind:预编译阶段,声明输出列的类型和名称
  • init:执行前阶段,打开游标或建立连接
  • apply:迭代阶段,每次填充一个数据块(最多 2048 行),返回 0 表示结束
  • initLocal(可选):多线程场景下的每线程初始化

实战:构建 mongo_query() 表函数

第一步:项目搭建

创建一个 Maven 项目,添加两个依赖:

<dependencies>
    <dependency>
        <groupId>org.duckdb</groupId>
        <artifactId>duckdb_jdbc</artifactId>
        <version>2.0.0-alpha</version>
    </dependency>
    <dependency>
        <groupId>org.mongodb</groupId>
        <artifactId>mongodb-driver-sync</artifactId>
        <version>5.1.0</version>
    </dependency>
</dependencies>

第二步:定义参数类

public class MongoQueryParameters {
    private String collectionName;
    private String queryJson;
    private String columns;
    private String hostname;
    private int port;
    private String database;
    private String username;
    private String password;

    // getters and setters...
}

第三步:实现表函数核心类

public class MongoQueryFunction implements DuckDBTableFunction {

    // ========== BIND: 声明输出 Schema ==========
    @Override
    public DuckDBTableFunctionBindData bind(DuckDBTableFunctionBindInfo info) throws Exception {
        String collectionName = info.getParameter(0).getString();
        String queryJson = info.getParameter(1).getString();
        String columnsJson = info.getNamedParameter("columns").getString();
        String hostname = info.getNamedParameter("hostname").getString();
        int port = info.getNamedParameter("port").getInt();
        String database = info.getNamedParameter("database").getString();
        String username = info.getNamedParameter("username").getString();
        String password = info.getNamedParameter("password").getString();

        // 解析 columns 参数,声明输出列
        List<String> columnNames = new ArrayList<>();
        for (BsonValue bv : BsonArray.parse(columnsJson)) {
            String name = bv.asString().getValue();
            columnNames.add(name);
            info.addResultColumn(name, String.class); // 全部声明为 String
        }

        // 建立 MongoDB 连接
        MongoClientSettings settings = MongoClientSettings.builder()
            .serverAddress(new ServerAddress(hostname, port))
            .applyToClusterSettings(builder -> 
                builder.hosts(List.of(new ServerAddress(hostname, port))))
            .build();
        MongoClient client = MongoClients.create(settings);
        MongoCollection<Document> collection = 
            client.getDatabase(database).getCollection(collectionName);

        Document query = Document.parse(queryJson);

        // 返回绑定数据,携带连接和查询信息
        return new MongoQueryBindData(client, collection, columnNames, query);
    }

    // ========== INIT: 打开游标 ==========
    @Override
    public DuckDBTableFunctionInitData init(DuckDBTableFunctionInitInfo info) throws Exception {
        info.setMaxThreads(1); // 本示例使用单线程
        MongoQueryBindData bindData = info.getBindData();
        FindIterable<Document> iter = bindData.collection.find(bindData.query);
        return new MongoQueryInitData(iter.cursor());
    }

    // ========== APPLY: 流式读取数据 ==========
    @Override
    public long apply(DuckDBTableFunctionCallInfo info, DuckDBDataChunkWriter output) throws Exception {
        MongoCursor<Document> cursor = info.getInitData().getResultCursor();
        long row = 0;
        
        for (; row < output.capacity() && cursor.hasNext(); row++) {
            Document doc = cursor.next();
            for (long col = 0; col < output.columnCount(); col++) {
                copyValueToVector(doc, output.vector(col), row, 
                    ((MongoQueryBindData)info.getBindData())
                        .getColumnNames().get((int)col));
            }
        }
        return row; // 返回 0 表示数据耗尽
    }
}

第四步:注册表函数

try (Connection conn = DriverManager.getConnection("jdbc:duckdb:")) {
    DuckDBFunctions.tableFunction()
        .withName("mongo_query")
        .withParameter(String.class)      // collection name
        .withParameter(String.class)      // Mongo filter (JSON)
        .withNamedParameter("columns", String.class)
        .withNamedParameter("hostname", String.class)
        .withNamedParameter("port", Integer.class)
        .withNamedParameter("database", String.class)
        .withNamedParameter("username", String.class)
        .withNamedParameter("password", String.class)
        .withFunction(new MongoQueryFunction())
        .register(conn);

    // 现在可以直接在 SQL 中使用!
    ResultSet rs = conn.executeQuery(
        "SELECT * FROM mongo_query('orders', " +
        "'{\"status\": \"shipped\"}', " +
        "columns='[\"customer_id\", \"amount\"]', " +
        "hostname='localhost', port=27017, database='app')"
    );
}

联邦查询实战:跨数据源 JOIN

表函数注册后,最强大的场景是异构数据源的联邦查询——远程 MongoDB 数据与本地 CSV 文件在同一 SQL 中 JOIN:

SELECT 
    c.region,
    COUNT(*) AS orders,
    SUM(o.amount::DECIMAL(10,2)) AS revenue
FROM mongo_query(
    'orders', 
    '{"status": "shipped"}',
    columns='["customer_id", "amount"]',
    database='app'
) AS o
JOIN 'customers/*.csv' AS c
ON c.customer_id = o.customer_id
GROUP BY c.region;

这条 SQL 做了什么?

  1. mongo_query() 从远程 MongoDB 流式读取已发货订单
  2. 'customers/*.csv' 读取本地客户 CSV 文件(支持 glob 通配)
  3. 在两数据源之间执行 JOIN,按地区聚合

没有数据导出步骤,没有中间表,纯粹的一条 SQL。


与传统方案的对比

特性纯 Java 表函数C++ 扩展导出+加载Trino 联邦查询
开发语言Java/Kotlin/ScalaC++任意SQL
构建复杂度低(Maven)高(CMake + DB API)
崩溃风险无(JVM 管理)Segfault
内存占用流式(~2048行/批次)可优化全量加载流式
谓词下推✅ 源端执行
多数据源 JOIN
部署方式普通 JAR共享库集群
适用规模单节点单节点单节点分布式

当前限制与注意事项

1. 仅支持 DuckDB Java 客户端

纯 Java 表函数不能作为 DuckDB 扩展加载(扩展是 Native 共享库)。这意味着表函数只能通过 DuckDB Java 客户端(JDBC/ODBC)使用,无法在 CLI 或 Python 中直接使用。

2. 生命周期管理需手动

当前版本中,bind()init() 返回的对象需要调用者自行管理生命周期。建议在对象上实现 AutoCloseable

public class MongoQueryInitData implements AutoCloseable {
    private final MongoCursor<Document> cursor;
    
    @Override
    public void close() {
        if (cursor != null) cursor.close();
    }
}

3. 复合类型暂不支持

当前向量 API 仅支持标量类型(STRING、INTEGER、DOUBLE 等)。STRUCT、LIST 等嵌套类型尚未支持,嵌套数据需要先展平或序列化为字符串。

4. 多线程执行需额外处理

示例使用 setMaxThreads(1) 限制单线程执行。如果需要多线程并行扫描,需要实现 initLocal 回调并为每个线程维护独立状态。


变现建议

方向 1:企业内部数据查询平台

痛点:企业数据分散在 Oracle、MongoDB、MySQL、CSV 文件等多个系统中,业务人员需要跨系统数据但缺乏技术能力。

方案:用纯 Java 表函数搭建一个内部数据查询平台:

  • 注册各数据源的表函数(Oracle JDBC、MongoDB、S3 Parquet 等)
  • 提供 SQL 查询界面(如 SQLPad 或自研 Web UI)
  • 业务人员用 SQL 跨系统查询,无需关心数据位置

变现模式

  • 一次性实施费:10,000-50,000 元
  • 年度维护费:20,000-100,000 元
  • 按查询配额收费:5-20 元/次高级查询

方向 2:数据分析 SaaS 插件市场

痛点:独立开发者需要快速接入各种数据源构建数据分析产品,但编写 C++ 扩展门槛太高。

方案:在 GitHub 上发布一系列纯 Java 表函数扩展(Salesforce、Stripe、Shopify API 等),每个作为独立 Maven 库:

  • 基础版免费开源
  • 企业版(带缓存、连接池、权限控制)收费订阅

变现模式

  • GitHub Sponsors:月入 500-5,000 元
  • 企业版订阅:$29-199/月/开发者
  • 技术咨询:按小时收费 $50-200

方向 3:数据中介服务(Data Broker)

痛点:中小企业有数据需求但无力自建数据管道,愿意为干净的分析数据付费。

方案:用 Java 表函数实时查询多个数据源,将结果打包为标准化数据产品:

  • 电商销售数据(Shopify + Stripe + MongoDB 聚合)
  • 社交媒体趋势数据(Twitter API + 存储分析)
  • 金融市场数据(券商 API + 新闻源)

变现模式

  • 数据订阅:¥99-999/月
  • API 调用计费:¥0.01-0.1/次
  • 定制数据产品:¥5,000-50,000/项目

总结

DuckDB 的纯 Java 表函数功能是嵌入式分析领域的重大突破。它让 Java 开发者能够以熟悉的工具(Maven、JDBC Driver)将任何数据源暴露为 SQL 表,与本地文件和其他数据源进行联邦查询,同时避免了 C++ 扩展的复杂性和风险。

随着 DuckDB v2.0 的发布和功能的完善,这个特性将成为 Java 生态中数据分析基础设施的重要组成部分。无论是企业内部的跨系统查询平台,还是独立开发者的数据 SaaS 产品,纯 Java 表函数都提供了前所未有的灵活性和开发效率。

📖 官方文档DuckDB Java Table Functions
📦 示例代码duckdb_mongo_example
💬 社区交流:DuckDB Discord #java 频道

📺 Watch video tutorials → Olap Studio YouTube

Subscribe for more DuckDB & AI automation tutorials

使用 Hugo 构建
主题 StackJimmy 设计

⚠️ 本站为独立社区项目,与 DuckDB 基金会及 DuckDB 官方项目无任何从属、背书或赞助关系。

"DuckDB" 是 DuckDB 基金会的注册商标,本站仅以事实描述方式使用该名称。

本站内容仅供教育与社区推广用途,不构成任何商业服务。