Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Java 客户端

用于通过协议与数据库服务器通信的 Java 客户端库。当前实现仅支持 HTTP interface。 该库提供了自有 API,用于向服务器发送请求,同时还提供了处理不同二进制数据格式 (RowBinary* & Native*) 的工具。

设置


<dependency>
    <groupId>com.clickhouse</groupId>
    <artifactId>client-v2</artifactId>
    <version>0.9.8</version>
</dependency>

初始化

Client 对象通过 com.clickhouse.client.api.Client.Builder#build() 初始化。每个客户端拥有独立的上下文,客户端之间不共享任何对象。 Builder 提供了多种配置方法,方便快速完成设置。

示例:

 Client client = new Client.Builder()
                .addEndpoint("https://clickhouse-cloud-instance:8443/")
                .setUsername(user)
                .setPassword(password)
                .build();

Client 实现了 AutoCloseable 接口,不再需要时应将其关闭。

身份验证

身份验证在初始化阶段按客户端进行配置。支持三种身份验证方法:密码、访问令牌以及 SSL 客户端证书。

通过密码进行身份验证需要调用 setUsername(String)setPassword(String) 来设置用户名和密码:

 Client client = new Client.Builder()
        .addEndpoint("https://clickhouse-cloud-instance:8443/")
        .setUsername(user)
        .setPassword(password)
        .build();

使用访问令牌进行身份验证需要调用 setAccessToken(String) 来设置访问令牌:

 Client client = new Client.Builder()
        .addEndpoint("https://clickhouse-cloud-instance:8443/")
        .setAccessToken(userAccessToken)
        .build();

通过 SSL 客户端证书进行身份验证,需要分别调用 setUsername(String)useSSLAuthentication(boolean)setClientCertificate(String)setClientKey(String) 来设置用户名、启用 SSL 身份验证、设置客户端证书和客户端密钥:

Client client = new Client.Builder()
        .useSSLAuthentication(true)
        .setUsername("some_user")
        .setClientCertificate("some_user.crt")
        .setClientKey("some_user.key")

配置

所有设置均通过实例方法 (即配置方法) 定义,使每个值的作用域和上下文一目了然。 主要配置参数在单一作用域 (客户端或操作) 内定义,彼此之间不会相互覆盖。

配置在客户端创建时定义。请参阅 com.clickhouse.client.api.Client.Builder

客户端配置

方法 参数 描述 默认值
addEndpoint(String endpoint) endpoint - URL 格式的服务器地址 将服务器端点添加到可用服务器列表。目前仅支持一个端点。 none none
addEndpoint(Protocol protocol, String host, int port, boolean secure) protocol - 连接协议
host - IP 或主机名
secure - 使用 HTTPS
将服务器端点添加到可用服务器列表。目前仅支持一个端点。 none none
enableConnectionPool(boolean enable) enable - 用于启用/禁用的标志 设置是否启用连接池。 true connection_pool_enabled
setMaxConnections(int maxConnections) maxConnections - 连接数 设置客户端可为每个服务器端点打开的连接数。 10 max_open_connections
setConnectionTTL(long timeout, ChronoUnit unit) timeout - 超时值
unit - 时间单位
设置连接的生存时间 (TTL),超过该时间后,连接将被视为非活动状态。 -1 connection_ttl
setKeepAliveTimeout(long timeout, ChronoUnit unit) timeout - 超时值
unit - 时间单位
设置 HTTP 连接的 Keep-Alive 超时时间。设为 0 可禁用 Keep-Alive。 - http_keep_alive_timeout
setConnectionReuseStrategy(ConnectionReuseStrategy strategy) strategy - LIFOFIFO 选择连接池使用的策略。 FIFO connection_reuse_strategy
setDefaultDatabase(String database) database - 数据库名称 设置默认数据库。 default database

客户端识别

查询日志中有两个字段用于标识发起请求的应用程序:client_namehttp_user_agent。Native TCP 协议使用 client_name 标识应用程序,HTTP 协议使用 http_user_agent 标识应用程序。Client builder 提供了 setClientName 方法,可为两种协议正确设置相应的值。 字段 http_user_agent 遵循 User-Agent 请求头的通用格式:application-name[/version] [(operating-system; architecture; ...)]。 这组值会在每个 layer 中重复出现:应用程序层、客户端库层、HTTP 客户端库层。由 setClientName 方法设置的内容在列表中排在最前面。

例如:

client.setClientName("my-app-01/1.0");

将产生以下 http_user_agent 值:

my-app-01/1.0 clickhouse-java-v2/0.9.6-SNAPSHOT (Linux; jvm:17.0.17) Apache-HttpClient/5.4.4

应用程序可以设置自定义 HTTP 请求头 User-Agent 来标识自身,但 clickhouse-java-v2/0.9.6-SNAPSHOT 部分会自动追加到该请求头的末尾。

操作标识

查询日志还有另外两个字段 query_idlog_comment,可用于标识操作并向查询日志添加额外信息。

query_id 是操作的唯一标识符。应用程序可通过调用 QuerySettings 类的 setQueryId 方法来设置该值。

QuerySettings querySettings = new QuerySettings();
querySettings.setQueryId("some-query-id");

log_comment 是一个可以添加到查询日志中的注释。应用程序可以通过调用 QuerySettings 类的 logComment 方法来设置此注释。

QuerySettings querySettings = new QuerySettings();
querySettings.logComment("some-comment");

服务器设置

服务器端配置可以在创建客户端时统一设置 (参见 BuilderserverSetting 方法) ,也可以在操作级别单独设置 (参见操作配置类的 serverSetting) 。

 try (Client client = new Client.Builder().addEndpoint(Protocol.HTTP, "localhost", mockServer.port(), false)
        .setUsername("default")
        .setPassword(ClickHouseServerForTest.getPassword())
        .compressClientRequest(true)

        // Client level
        .serverSetting("max_threads", "10")
        .serverSetting("async_insert", "1")
        .serverSetting("roles", Arrays.asList("role1", "role2"))

        .build()) {

	// Operation level
	QuerySettings querySettings = new QuerySettings();
	querySettings.serverSetting("session_timezone", "Europe/Zurich");

	...
}

⚠️ 当通过 setOption 方法 (Client.Builder 或操作设置类) 设置选项时,server settings 名称应加上 clickhouse_setting_ 前缀。此时 com.clickhouse.client.api.ClientConfigProperties#serverSetting() 会很有帮助。

自定义 HTTP 请求头

可以为所有操作 (客户端级别) 或单个操作 (操作级别) 设置自定义 HTTP 请求头。


QuerySettings settings = new QuerySettings()
    .httpHeader(HttpHeaders.REFERER, clientReferer)
    .setQueryId(qId);

通过 setOption 方法 (Client.Builder 或操作设置类) 设置选项时,自定义请求头名称应以 http_header_ 为前缀。在这种情况下,com.clickhouse.client.api.ClientConfigProperties#httpHeader() 方法会很有帮助。

常用定义

ClickHouseFormat

支持格式的枚举类型,包含 ClickHouse 支持的所有格式。

  • raw - 用户应对原始数据进行转码
  • full - 客户端可以自行转码数据,并接收原始数据流
  • - - ClickHouse 不支持此格式的该操作

此客户端版本支持:

格式 输入 输出
TabSeparated raw raw
TabSeparatedRaw raw raw
TabSeparatedWithNames raw raw
TabSeparatedWithNamesAndTypes raw raw
TabSeparatedRawWithNames raw raw
TabSeparatedRawWithNamesAndTypes raw raw
Template raw raw
TemplateIgnoreSpaces raw -
CSV raw raw
CSVWithNames raw raw
CSVWithNamesAndTypes raw raw
CustomSeparated raw raw
CustomSeparatedWithNames raw raw
CustomSeparatedWithNamesAndTypes raw raw
SQLInsert - raw
Values raw raw
Vertical - raw
JSON raw raw
JSONAsString raw -
JSONAsObject raw -
JSONStrings raw raw
JSONColumns raw raw
JSONColumnsWithMetadata raw raw
JSONCompact raw raw
JSONCompactStrings - raw
JSONCompactColumns raw raw
JSONEachRow raw raw
PrettyJSONEachRow - raw
JSONEachRowWithProgress - raw
JSONStringsEachRow raw raw
JSONStringsEachRowWithProgress - raw
JSONCompactEachRow raw raw
JSONCompactEachRowWithNames raw raw
JSONCompactEachRowWithNamesAndTypes raw raw
JSONCompactStringsEachRow raw raw
JSONCompactStringsEachRowWithNames raw raw
JSONCompactStringsEachRowWithNamesAndTypes raw raw
JSONObjectEachRow raw raw
BSONEachRow raw raw
TSKV raw raw
Pretty - raw
PrettyNoEscapes - raw
PrettyMonoBlock - raw
PrettyNoEscapesMonoBlock - raw
PrettyCompact - raw
PrettyCompactNoEscapes - raw
PrettyCompactMonoBlock - raw
PrettyCompactNoEscapesMonoBlock - raw
PrettySpace - raw
PrettySpaceNoEscapes - raw
PrettySpaceMonoBlock - raw
PrettySpaceNoEscapesMonoBlock - raw
Prometheus - raw
Protobuf raw raw
ProtobufSingle raw raw
ProtobufList raw raw
Avro raw raw
AvroConfluent raw -
Parquet raw raw
ParquetMetadata raw -
Arrow raw raw
ArrowStream raw raw
ORC raw raw
One raw -
Npy raw raw
RowBinary full full
RowBinaryWithNames full full
RowBinaryWithNamesAndTypes full full
RowBinaryWithDefaults full -
Native full raw
Null - raw
XML - raw
CapnProto raw raw
LineAsString raw raw
Regexp raw -
RawBLOB raw raw
MsgPack raw raw
MySQLDump raw -
DWARF raw -
Markdown - raw
Form raw -

插入 API

insert(String tableName, InputStream data, ClickHouseFormat format)

以指定格式接受字节 InputStream 形式的数据。data 须以 format 指定的格式进行编码。

签名

CompletableFuture<InsertResponse> insert(String tableName, InputStream data, ClickHouseFormat format, InsertSettings settings)
CompletableFuture<InsertResponse> insert(String tableName, InputStream data, ClickHouseFormat format)

参数

tableName - 目标表名称。

data - 已编码数据的输入流。

format - 数据编码所使用的格式。

settings - 请求设置。

返回值

InsertResponse 类型的 Future——操作结果及附加信息,例如服务器端指标。

示例

try (InputStream dataStream = getDataStream()) {
    try (InsertResponse response = client.insert(TABLE_NAME, dataStream, ClickHouseFormat.JSONEachRow,
            insertSettings).get(3, TimeUnit.SECONDS)) {

        log.info("Insert finished: {} rows written", response.getMetrics().getMetric(ServerMetrics.NUM_ROWS_WRITTEN).getLong());
    } catch (Exception e) {
        log.error("Failed to write JSONEachRow data", e);
        throw new RuntimeException(e);
    }
}

insert(String tableName, List<?> data, InsertSettings settings)

向数据库发送写请求。对象列表将被转换为高效格式后发送至服务器。列表项的类需要预先通过 register(Class, TableSchema) 方法注册。

签名

client.insert(String tableName, List<?> data, InsertSettings settings)
client.insert(String tableName, List<?> data)

参数

tableName - 目标表的名称。

data - DTO (数据传输对象) 对象集合。

settings - 请求设置。

返回值

InsertResponse 类型的 Future——包含操作结果及服务端指标等附加信息。

示例

// Important step (done once) - register class to pre-compile object serializer according to the table schema.
client.register(ArticleViewEvent.class, client.getTableSchema(TABLE_NAME));

List<ArticleViewEvent> events = loadBatch();

try (InsertResponse response = client.insert(TABLE_NAME, events).get()) {
    // handle response, then it will be closed and connection that served request will be released.
}

InsertSettings

插入操作的配置选项。

配置方法

方法 描述
setQueryId(String queryId) 设置要分配给该操作的查询 ID。默认值:null
setDeduplicationToken(String token) 设置去重令牌。该令牌会发送到服务器,可用于标识查询。默认值:null
setInputStreamCopyBufferSize(int size) 复制缓冲区大小。该缓冲区在写入操作期间用于将数据从用户提供的输入流复制到输出流。默认值:8196
serverSetting(String name, String value) 为某个操作设置单个服务器级设置项。
serverSetting(String name, Collection values) 为某项操作设置一个具有多个值的服务器设置。集合中的项应为 String 值。
setDBRoles(Collection dbRoles) 设置在执行操作前生效的 DB 角色。集合中的项应为 String 值。
setOption(String option, Object value) 以原始格式设置一个配置选项。这不是服务器级设置项。

InsertResponse

保存插入操作结果的响应对象。仅在客户端收到服务器响应时可用。

方法 描述
OperationMetrics getMetrics() 返回包含操作指标的对象。
String getQueryId() 返回该操作的查询 ID,该 ID 由应用程序通过操作设置指定,或由服务器分配。

查询 API

query(String sqlQuery)

按原样发送 sqlQuery。响应格式由查询设置决定。QueryResponse 将保存对响应流的引用,该流应由支持相应格式的读取器 (reader) 读取。

签名

CompletableFuture<QueryResponse> query(String sqlQuery, QuerySettings settings)
CompletableFuture<QueryResponse> query(String sqlQuery)

参数

sqlQuery - 单条 SQL 语句。该查询会按原样发送到服务器。

settings - 请求设置。

返回值

QueryResponse 类型的 Future——包含结果数据集及服务端指标等附加信息。消费完数据集后,应关闭 Response 对象。

示例

final String sql = "select * from " + TABLE_NAME + " where title <> '' limit 10";

// Default format is RowBinaryWithNamesAndTypesFormatReader so reader have all information about columns
try (QueryResponse response = client.query(sql).get(3, TimeUnit.SECONDS);) {

    // Create a reader to access the data in a convenient way
    ClickHouseBinaryFormatReader reader = client.newBinaryFormatReader(response);

    while (reader.hasNext()) {
        reader.next(); // Read the next record from stream and parse it

        // get values
        double id = reader.getDouble("id");
        String title = reader.getString("title");
        String url = reader.getString("url");

        // collecting data
    }
} catch (Exception e) {
    log.error("Failed to read data", e);
}

// put business logic outside of the reading block to release http connection asap.

query(String sqlQuery, Map<String, Object> queryParams, QuerySettings settings)

按原样发送 sqlQuery,同时发送查询参数,以便服务器编译 SQL 表达式。

签名

CompletableFuture<QueryResponse> query(String sqlQuery, Map<String, Object> queryParams, QuerySettings settings)

参数

sqlQuery - 带占位符 {} 的 SQL 表达式。

queryParams - 用于在服务器端补全 SQL 表达式的变量映射。

settings - 请求设置。

返回值

QueryResponse 类型的 Future——包含结果数据集及服务端指标等附加信息。消费完数据集后,应关闭 Response 对象。

示例


// define parameters. They will be sent to the server along with the request.
Map<String, Object> queryParams = new HashMap<>();
queryParams.put("param1", 2);

try (QueryResponse response =
        client.query("SELECT * FROM " + table + " WHERE col1 >= {param1:UInt32}", queryParams, new QuerySettings()).get()) {

    // Create a reader to access the data in a convenient way
    ClickHouseBinaryFormatReader reader = client.newBinaryFormatReader(response);

    while (reader.hasNext()) {
        reader.next(); // Read the next record from stream and parse it

        // reading data
    }

} catch (Exception e) {
    log.error("Failed to read data", e);
}

queryAll(String sqlQuery)

RowBinaryWithNamesAndTypes 格式查询数据,并以集合形式返回结果。读取性能与使用 reader 时相同,但需要占用更多内存来存储完整的数据集。

签名

List<GenericRecord> queryAll(String sqlQuery)

参数

sqlQuery - 用于从服务器查询数据的 SQL 表达式。

返回值

完整数据集,由一组 GenericRecord 对象列表表示,以行方式访问结果数据。

示例

try {
    log.info("Reading whole table and process record by record");
    final String sql = "select * from " + TABLE_NAME + " where title <> ''";

    // Read whole result set and process it record by record
    client.queryAll(sql).forEach(row -> {
        double id = row.getDouble("id");
        String title = row.getString("title");
        String url = row.getString("url");

        log.info("id: {}, title: {}, url: {}", id, title, url);
    });
} catch (Exception e) {
    log.error("Failed to read data", e);
}

QuerySettings

查询操作的配置选项。

配置方法

方法 描述
setQueryId(String queryId) 设置将分配给该操作的查询 ID。
setFormat(ClickHouseFormat format) 设置响应格式。完整列表请参见 RowBinaryWithNamesAndTypes
setMaxExecutionTime(Integer maxExecutionTime) 设置该操作在服务器端的执行时间。不会影响读取超时。
waitEndOfQuery(Boolean waitEndOfQuery) 请求服务器在发送响应前等待查询结束。
setUseServerTimeZone(Boolean useServerTimeZone) 将使用服务器时区 (参见客户端配置) 来解析操作结果中的日期/时间类型。默认为 false
setUseTimeZone(String timeZone) 请求服务器使用 timeZone 进行时间转换。请参见 session_timezone
serverSetting(String name, String value) 为该操作设置单个服务器级设置。
serverSetting(String name, Collection values) 为某项操作设置一个具有多个值的服务器设置。集合中的项应为 String 值。
setDBRoles(Collection dbRoles) 设置在执行操作前生效的 DB 角色。集合中的项应为 String 值。
setOption(String option, Object value) 以原始格式设置一个配置选项。这不是服务器级设置项。

QueryResponse

保存查询执行结果的响应对象。仅当客户端收到来自服务器的响应时,该对象才可用。

方法 描述
ClickHouseFormat getFormat() 返回响应中数据采用的编码格式。
InputStream getInputStream() 返回采用指定格式的未压缩数据字节流。
OperationMetrics getMetrics() 返回包含操作指标的对象。
String getQueryId() 返回该操作的查询 ID,该 ID 由应用程序分配 (通过操作设置) 或由服务器分配。
TimeZone getTimeZone() 返回处理响应中的 Date/DateTime 类型时应使用的时区。

示例

  • 示例代码可在仓库中找到
  • 参考 Spring Service 的实现

通用 API

getTableSchema(String table)

获取 table 的 schema。

签名

TableSchema getTableSchema(String table)
TableSchema getTableSchema(String table, String database)

参数

table - 需要拉取 schema 数据的表名。

database - 目标表所在的数据库。

返回值

返回一个包含表列列表的 TableSchema 对象。

getTableSchemaFromQuery(String sql)

从 SQL 语句中拉取 schema。

签名

TableSchema getTableSchemaFromQuery(String sql)

参数

sql - 需要返回其 schema 的 "SELECT" SQL 语句。

返回值

返回一个 TableSchema 对象,其列与 sql 表达式相匹配。

TableSchema

register(Class<?> clazz, TableSchema schema)

为 Java 类编译序列化与反序列化 layer,以便通过 schema 进行数据读写。该 method 将为 getter/setter 对及其对应列创建序列化器和反序列化器。 列的匹配通过从方法名中提取列名来实现。例如,getFirstName 对应的列为 first_namefirstname

签名

void register(Class<?> clazz, TableSchema schema)

参数

clazz - 表示用于读取/写入数据的 POJO 的类。

schema - 用于匹配 POJO 属性的 schema。

示例

client.register(ArticleViewEvent.class, client.getTableSchema(TABLE_NAME));

使用示例

完整的示例代码存放在代码仓库的 'example` 文件夹中:

读取数据

读取数据有两种常见方式:

  • 返回底层 QueryResponse 对象的 query() 方法,该对象包含带有数据的 InputStream。通常会结合 ClickHouseBinaryFormatReader 用于流式读取,但 也可搭配任何其他自定义 reader 实现使用。QueryResponse 还提供对结果集元数据和指标的访问。
  • queryAll() 方法,并使用 GenericRecord 以便更方便地访问行。在这种情况下,整个结果集都会加载到内存中。
  • 返回 com.clickhouse.client.api.query.RecordsqueryRecords() 方法——它是一个用于遍历 GenericRecord 对象的迭代器。该方法采用流式方式 (不会将任何数据加载到内存中) ,并通过 GenericRecord 访问数据。

注意: 流式方法要求快速读取,否则可能导致服务器写入超时,因为数据是直接从网络流中读取的。

读取数组

ClickHouseBinaryFormatReader 方法

  • getList(...) - 将任意 Array(...) 读取为 List<T>。这是灵活类型读取的默认推荐选择。支持嵌套数组。
  • getByteArray(...), getShortArray(...), getIntArray(...), getLongArray(...), getFloatArray(...), getDoubleArray(...), getBooleanArray(...) - 最适合用于包含与基本类型兼容值的一维数组。
  • getStringArray(...) - 用于 Array(String) (以及以名称形式表示的枚举值) 。
  • getObjectArray(...) - 适用于任意 Array(...) 元素类型的通用选项,包括嵌套数组。可用于读取包含 Nullable 值和嵌套数组的数组。

所有方法均支持基于索引和基于名称的重载。索引从 1 开始计数。基于索引的方式可直接访问列。 基于名称的方法每次调用时均需执行索引查找。

try (QueryResponse response = client.query("SELECT * FROM my_table").get()) {
    ClickHouseBinaryFormatReader reader = client.newBinaryFormatReader(response);
    while (reader.next() != null) {

        Object[] uint64 = reader.getObjectArray("uint64_arr"); // Array(UInt64) -> BigInteger[]
        Object[] arr2d = reader.getObjectArray("arr2d");       // Array(Array(Int64)) -> Object[]

        // nested arrays are returned as nested Object[]:
        Object[] firstInner = (Object[]) arr2d[0];
        Long firstValue = (Long) firstInner[0];
    }
}

GenericRecord 方法

  • getList(...) - 将任意 Array(...) 读取为 List<T>。这是灵活类型读取的默认推荐选择。支持嵌套数组。
  • getByteArray(...), getShortArray(...), getIntArray(...), getLongArray(...), getFloatArray(...), getDoubleArray(...), getBooleanArray(...) - 最适合用于包含与基本类型兼容值的一维数组。
  • getStringArray(...) - 用于 Array(String) (以及以名称形式表示的枚举值) 。
  • getObjectArray(...) - 适用于任意 Array(...) 元素类型的通用选项,包括嵌套数组。可用于读取包含 Nullable 值和嵌套数组的数组。

所有方法均支持基于索引和基于名称的重载。索引从 1 开始计数。基于索引的方式可直接访问列。 基于名称的方法每次调用时均需执行索引查找。

try (QueryResponse response = client.query("SELECT * FROM my_table").get()) {
    List<GenericRecord> rows = client.queryAll(
        "SELECT int_arr, arr2d_nullable FROM test_arrays ORDER BY id");

    for (GenericRecord row : rows) {
        Object[] intArr = row.getObjectArray("int_arr");                 // Array(Int32) -> Integer[]
        Object[] arr2d = row.getObjectArray("arr2d_nullable");           // Array(Array(Nullable(Int32)))

        Object[] inner = (Object[]) arr2d[0];
        Object maybeNull = inner[1]; // may be null
    }
}

迁移指南

旧版客户端 (V1) 以 com.clickhouse.client.ClickHouseClient#builder 作为起点。新版客户端 (V2) 采用类似的模式,使用 com.clickhouse.client.api.Client.Builder。主要区别如下:

  • 未使用服务加载器来获取具体实现。com.clickhouse.client.api.Client 是一个门面类,未来可适配各种实现。
  • 配置来源更少:一类提供给构建器,另一类在操作设置 (QuerySettingsInsertSettings) 中。先前版本会为每个节点分别配置,并且在某些情况下会加载 环境变量。

配置参数匹配

V1 中有 3 个与配置相关的枚举类:

  • com.clickhouse.client.config.ClickHouseDefaults - 在大多数用例中都应设置的配置参数,例如 USERPASSWORD
  • com.clickhouse.client.config.ClickHouseClientOption - 客户端特有的配置参数,例如 HEALTH_CHECK_INTERVAL
  • com.clickhouse.client.http.config.ClickHouseHttpOption - HTTP 接口专用的配置参数,例如 RECEIVE_QUERY_PROGRESS

这些设计的初衷是对参数进行分组并实现清晰的隔离。然而在某些情况下,这反而造成了混淆 (例如,com.clickhouse.client.config.ClickHouseDefaults#ASYNCcom.clickhouse.client.config.ClickHouseClientOption#ASYNC 之间究竟有何区别?) 。新的 V2 客户端以 com.clickhouse.client.api.Client.Builder 作为所有可用客户端配置选项的统一字典。所有配置参数名称均列于 com.clickhouse.client.api.ClientConfigProperties 中。

下表列出了新客户端中支持的旧选项及其新含义。

图例: ✔ = 支持,✗ = 已丢弃

V1 配置 V2 构建器方法 备注
ClickHouseDefaults#HOST Client.Builder#addEndpoint
ClickHouseDefaults#PROTOCOL V2 仅支持 HTTP
ClickHouseDefaults#DATABASE
ClickHouseClientOption#DATABASE
Client.Builder#setDefaultDatabase
ClickHouseDefaults#USER Client.Builder#setUsername
ClickHouseDefaults#PASSWORD Client.Builder#setPassword
ClickHouseClientOption#CONNECTION_TIMEOUT Client.Builder#setConnectTimeout
ClickHouseClientOption#CONNECTION_TTL Client.Builder#setConnectionTTL
ClickHouseHttpOption#MAX_OPEN_CONNECTIONS Client.Builder#setMaxConnections
ClickHouseHttpOption#KEEP_ALIVE
ClickHouseHttpOption#KEEP_ALIVE_TIMEOUT
Client.Builder#setKeepAliveTimeout
ClickHouseHttpOption#CONNECTION_REUSE_STRATEGY Client.Builder#setConnectionReuseStrategy
ClickHouseHttpOption#USE_BASIC_AUTHENTICATION Client.Builder#useHTTPBasicAuth

主要差异

  • Client V2 使用更少的专有类,以提高可移植性。例如,V2 可配合 java.io.InputStream 的任何实现来 将数据写入服务器。
  • Client V2 的 async 设置默认值为 off。这意味着不会额外创建线程,且应用程序对客户端有更强的控制能力。对于绝大多数用例,此设置都应保持为 off。启用 async 会为每个请求创建一个单独的线程。只有在使用由应用程序控制的 executor 时,这样做才有意义 (请参见 com.clickhouse.client.api.Client.Builder#setSharedOperationExecutor)

写入数据

  • 可使用 java.io.InputStream 的任何实现。支持 V1 com.clickhouse.data.ClickHouseInputStream,但不建议使用。
  • 一旦检测到输入流结束,就会进行相应处理。在此之前,应关闭请求的输出流。

V1 插入 TSV 格式的数据。

InputStream inData = getInData();
ClickHouseRequest.Mutation request = client.read(server)
        .write()
        .table(tableName)
        .format(ClickHouseFormat.TSV);
ClickHouseConfig config = request.getConfig();
CompletableFuture<ClickHouseResponse> future;
try (ClickHousePipedOutputStream requestBody = ClickHouseDataStreamFactory.getInstance()
        .createPipedOutputStream(config)) {
    // start the worker thread which transfer data from the input into ClickHouse
    future = request.data(requestBody.getInputStream()).execute();

    // Copy data from inData stream to requestBody stream

    // We need to close the stream before getting a response
    requestBody.close();

    try (ClickHouseResponse response = future.get()) {
        ClickHouseResponseSummary summary = response.getSummary();
        Assert.assertEquals(summary.getWrittenRows(), numRows, "Num of written rows");
    }
}

V2 插入 TSV 格式数据。

InputStream inData = getInData();
InsertSettings settings = new InsertSettings().setInputStreamCopyBufferSize(8198 * 2); // set copy buffer size
try (InsertResponse response = client.insert(tableName, inData, ClickHouseFormat.TSV, settings).get(30, TimeUnit.SECONDS)) {

  // Insert is complete at this point

} catch (Exception e) {
 // Handle exception
}
  • 只需调用一个方法。无需另外创建请求对象。
  • 请求体流会在所有数据复制完毕后自动关闭。
  • 现已提供新的底层 API:com.clickhouse.client.api.Client#insert(java.lang.String, java.util.List<java.lang.String>, com.clickhouse.client.api.DataStreamWriter, com.clickhouse.data.ClickHouseFormat, com.clickhouse.client.api.insert.InsertSettings)com.clickhouse.client.api.DataStreamWriter 用于实现自定义的数据写入逻辑。例如,从 队列中读取数据。

读取数据

  • 默认使用 RowBinaryWithNamesAndTypes 格式读取数据。当前在需要进行数据绑定时,仅支持此格式。
  • 可以使用 List<GenericRecord> com.clickhouse.client.api.Client#queryAll(java.lang.String) 方法,将数据作为记录集合读取。它会将数据读入内存并释放连接,无需额外处理。GenericRecord 提供对数据的访问,并支持一些类型转换。
Collection<GenericRecord> records = client.queryAll("SELECT * FROM table");
for (GenericRecord record : records) {
    int rowId = record.getInteger("rowID");
    String name = record.getString("name");
    LocalDateTime ts = record.getLocalDateTime("ts");
}

用于通过协议与数据库服务器通信的 Java 客户端库。当前实现仅支持 HTTP 接口。该库提供自有 API,用于向服务器发送请求。

设置

<!-- https://mvnrepository.com/artifact/com.clickhouse/clickhouse-http-client -->
<dependency>
    <groupId>com.clickhouse</groupId>
    <artifactId>clickhouse-http-client</artifactId>
    <version>0.7.2</version>
</dependency>

自版本 0.5.0 起,驱动程序引入了新的客户端 HTTP 库,需将其作为依赖项添加。

<!-- https://mvnrepository.com/artifact/org.apache.httpcomponents.client5/httpclient5 -->
<dependency>
    <groupId>org.apache.httpcomponents.client5</groupId>
    <artifactId>httpclient5</artifactId>
    <version>5.3.1</version>
</dependency>

初始化

连接 URL 格式:protocol://host[:port][/database][?param[=value][&param[=value]][#tag[,tag]],例如:

  • http://localhost:8443?ssl=true&sslmode=NONE
  • https://(https://explorer@play.clickhouse.com:443

连接到单个节点:

ClickHouseNode server = ClickHouseNode.of("http://localhost:8123/default?compress=0");

连接到包含多个节点的集群:

ClickHouseNodes servers = ClickHouseNodes.of(
    "jdbc:ch:http://server1.domain,server2.domain,server3.domain/my_db"
    + "?load_balancing_policy=random&health_check_interval=5000&failover=2");

查询 API

try (ClickHouseClient client = ClickHouseClient.newInstance(ClickHouseProtocol.HTTP);
     ClickHouseResponse response = client.read(servers)
        .format(ClickHouseFormat.RowBinaryWithNamesAndTypes)
        .query("select * from numbers limit :limit")
        .params(1000)
        .executeAndWait()) {
            ClickHouseResponseSummary summary = response.getSummary();
            long totalRows = summary.getTotalRowsToRead();
}

流式查询 API

try (ClickHouseClient client = ClickHouseClient.newInstance(ClickHouseProtocol.HTTP);
     ClickHouseResponse response = client.read(servers)
        .format(ClickHouseFormat.RowBinaryWithNamesAndTypes)
        .query("select * from numbers limit :limit")
        .params(1000)
        .executeAndWait()) {
            for (ClickHouseRecord r : response.records()) {
            int num = r.getValue(0).asInteger();
            // 类型转换
            String str = r.getValue(0).asString();
            LocalDate date = r.getValue(0).asDate();
        }
}

请参阅代码库中的完整代码示例

插入 API

try (ClickHouseClient client = ClickHouseClient.newInstance(ClickHouseProtocol.HTTP);
     ClickHouseResponse response = client.read(servers).write()
        .format(ClickHouseFormat.RowBinaryWithNamesAndTypes)
        .query("insert into my_table select c2, c3 from input('c1 UInt8, c2 String, c3 Int32')")
        .data(myInputStream) // `myInputStream` 是 RowBinary 格式的数据源
        .executeAndWait()) {
            ClickHouseResponseSummary summary = response.getSummary();
            summary.getWrittenRows();
}

请参阅代码库中的完整代码示例

RowBinary 编码

RowBinary 格式详见其页面

这里有一个代码示例

功能特性

压缩

客户端默认使用 LZ4 压缩,需要以下依赖项:

<!-- https://mvnrepository.com/artifact/org.lz4/lz4-java -->
<dependency>
    <groupId>org.lz4</groupId>
    <artifactId>lz4-java</artifactId>
    <version>1.8.0</version>
</dependency>

您也可以通过在连接 URL 中设置 compress_algorithm=gzip 来改用 gzip。

或者,您也可以通过以下几种方式禁用压缩。

  1. 在 connection URL 中将 compress=0 设为禁用:http://localhost:8123/default?compress=0
  2. 在客户端配置中禁用:
ClickHouseClient client = ClickHouseClient.builder()
   .config(new ClickHouseConfig(Map.of(ClickHouseClientOption.COMPRESS, false)))
   .nodeSelector(ClickHouseNodeSelector.of(ClickHouseProtocol.HTTP))
   .build();

请参阅压缩文档,了解各种压缩选项的详细信息。

多个查询

在同一会话中,在工作线程内依次执行多个查询:

CompletableFuture<List<ClickHouseResponseSummary>> future = ClickHouseClient.send(servers.apply(servers.getNodeSelector()),
    "create database if not exists my_base",
    "use my_base",
    "create table if not exists test_table(s String) engine=Memory",
    "insert into test_table values('1')('2')('3')",
    "select * from test_table limit 1",
    "truncate table test_table",
    "drop table if exists test_table");
List<ClickHouseResponseSummary> results = future.get();

命名参数

您可以按名称传递参数,而无需依赖其在参数列表中的位置。此功能可通过 params 函数实现。

try (ClickHouseClient client = ClickHouseClient.newInstance(ClickHouseProtocol.HTTP);
     ClickHouseResponse response = client.read(servers)
        .format(ClickHouseFormat.RowBinaryWithNamesAndTypes)
        .query("select * from my_table where name=:name limit :limit")
        .params("Ben", 1000)
        .executeAndWait()) {
            //...
        }
}

节点发现

Java client 提供了自动发现 ClickHouse 节点的功能。自动发现默认处于禁用状态。如需手动启用,请将 auto_discovery 设置为 true

properties.setProperty("auto_discovery", "true");

或在连接 URL 中:

jdbc:ch://my-server/system?auto_discovery=true

如果启用了自动发现功能,则无需在连接 URL 中指定所有 ClickHouse 节点。URL 中指定的节点将被视为种子节点,Java 客户端将自动从系统表和/或 clickhouse-keeper 或 zookeeper 中发现更多节点。

以下选项用于配置自动发现功能:

属性 默认值 描述
auto_discovery false 客户端是否应通过系统表和/或 clickhouse-keeper/zookeeper 发现更多节点。
node_discovery_interval 0 节点发现间隔 (毫秒) ;值为零或负数表示仅发现一次。
node_discovery_limit 100 一次最多可发现的节点数;值为零或负数表示不受限制。

负载均衡

Java 客户端根据负载均衡策略选择 ClickHouse 节点来发送请求。通常,负载均衡策略负责以下事项:

  1. 从托管节点列表中获取节点。
  2. 管理节点状态。
  3. 可选择为节点发现调度一个后台进程 (如果已启用自动发现) ,并执行健康检查。

以下是配置负载均衡的选项列表:

配置项 默认值 说明
load_balancing_policy "" 负载均衡策略可以是以下之一:
  • firstAlive - 请求会发送到托管节点列表中的第一个健康节点
  • random - 请求会发送到托管节点列表中的任意一个节点
  • roundRobin - 请求会依次发送到托管节点列表中的各个节点。
  • 实现 ClickHouseLoadBalancingPolicy 的完全限定类名 - 自定义负载均衡策略
  • 如果未指定,则请求会发送到托管节点列表中的第一个节点
    load_balancing_tags "" 用于筛选节点的负载均衡标签。请求只会发送到带有指定标签的节点
    health_check_interval 0 健康检查间隔,单位为毫秒;零值或负值表示一次性。
    health_check_method ClickHouseHealthCheckMethod.SELECT_ONE 健康检查方法。可以是以下任一项:
  • ClickHouseHealthCheckMethod.SELECT_ONE - 使用 select 1 查询进行检查
  • ClickHouseHealthCheckMethod.PING - 协议专用检查,通常更快
  • node_check_interval 0 节点检查间隔 (以毫秒为单位) ,负数按零处理。如果距离上次检查已超过指定时间,则会检查节点状态。
    health_check_intervalnode_check_interval 的区别在于,health_check_interval 选项会调度后台任务,用于检查节点列表 (全部节点或故障节点) 的状态,而 node_check_interval 则指定某个特定节点距离上次检查需要经过的时间
    check_all_nodes false 是否对所有节点执行健康检查,还是仅对故障节点执行健康检查。

    故障转移与重试

    Java client 提供了配置选项,用于为失败查询设置故障转移和重试行为:

    属性 默认值 说明
    故障转移 0 单个请求允许发生故障转移的最大次数。值为 0 或负数表示不进行故障转移。发生故障转移时,会根据负载均衡策略将失败的请求发送到其他节点,以从故障中恢复。
    重试 0 请求可重试的最大次数。零或负值表示不重试。只有在 ClickHouse server 返回 NETWORK_ERROR 错误码时,重试才会将请求发送到同一节点
    repeat_on_session_lock true 当会话被锁定时,是否持续重试执行直到超时 (根据 session_timeoutconnect_timeout) 。如果 ClickHouse 服务器返回 SESSION_IS_LOCKED 错误代码,则会重试失败的请求

    添加自定义 HTTP 请求头

    Java client 支持 HTTP/S 传输层,可用于向请求中添加自定义 HTTP 请求头。 请使用 custom_http_headers 属性,多个请求头之间以 , 分隔,请求头的键/值之间以 = 分隔。

    Java Client 支持

    options.put("custom_http_headers", "X-ClickHouse-Quota=test, X-ClickHouse-Test=test");

    JDBC 驱动

    properties.setProperty("custom_http_headers", "X-ClickHouse-Quota=test, X-ClickHouse-Test=test");
    Navigation