一、参数介绍

先看官网介绍关于数据insert分布式表的参数:

参数组1

  • distributed_foreground_insert
    默认值:0
    启用或禁用同步数据插入到 Distributed 表中。
    默认情况下,当将数据插入到 Distributed 表中时,ClickHouse 服务器以后台模式将数据发送到集群节点。当 distributed_foreground_insert=1 时,数据将同步处理,并且只有在所有数据都保存在所有分片上(如果 internal_replication 为 true,则每个分片至少一个副本)后,INSERT 操作才会成功。
    可能的值:
    0 — 数据以后台模式插入。
    1 — 数据以同步模式插入。

  • insert_distributed_sync
    参数介绍跟distribued_foreground_insert一样,新版本的参数已经修改成distributed_foreground_insert,老版本的参数还是insert_distributed_sync。

客户端执行insert的参数

参数组2

  • async_insert:是否异步执行insert语句
    默认值: 0
    如果为 true,则来自 INSERT 查询的数据存储在队列中,稍后在后台刷新到表中。如果 wait_for_async_insert 为 false,则 INSERT 查询几乎立即处理,否则客户端将等待直到数据刷新到表中
  • wait_for_async_insert:是否等待异步执行语句返回
    默认值: 1
    如果为 true,则等待异步插入的处理完成。
    注意:因为默认情况下async_insert参数和wait_for_async_insert参数是按照同步的方式进行配置,这里不单独进行设置。

区别

insert_distributed_sync和async_insert两个用于控制插入行为的参数,但它们在应用场景和机制上有重要区别。
先说结论:
在这里插入图片描述

insert_distributed_sync
作用:

  • 当向分布式表(Distributed table)插入数据时,控制插入行为是同步还是异步。
  • 默认情况下,insert_distributed_sync=0,即异步插入。这意味着数据先写入本地文件系统,然后再在后台异步地发送到集群分片。
  • 设置为insert_distributed_sync=1时,插入操作将同步进行,即客户端会等待数据成功写入所有分片(或达到超时时间)后才返回。

重要性和影响:

  • 数据一致性:同步插入可以确保在插入操作返回时,数据已经写入所有分片,从而保证强一致性。而在异步模式下,如果在后台发送数据之前发生故障,可能会导致数据丢失。
  • 性能:同步插入会降低写入吞吐量,因为需要等待网络传输和所有分片的写入完成。异步插入则具有更高的吞吐量,因为客户端不需要等待后台操作。
  • 超时控制:当使用同步插入时,可以通过insert_distributed_timeout参数设置等待超时时间,默认为180秒。如果超时,插入操作会失败。

使用场景:

  • 当需要确保数据已经安全写入所有分片时,例如在关键业务数据插入时,可以使用同步插入。

async_insert
作用:

  • 控制是否启用异步插入,该参数针对的是插入请求的处理方式,而不是分布式表的分片数据同步。
  • 默认情况下,async_insert=0,即每个插入请求都会立即被处理
  • 设置为async_insert=1时,ClickHouse会将插入请求排队,并定期批量处理它们。这可以减少小批量插入的频率,提高吞吐量。

重要性和影响:

  • 吞吐量:异步插入可以将多个插入请求合并成一批,减少IO操作和表锁竞争,从而提高吞吐量。
  • 延迟:异步插入会增加数据可见的延迟,因为数据不会立即被写入表中。
  • 可靠性:异步插入存在一定的风险,如果在数据被刷新到磁盘之前服务器发生故障,可能会丢失部分数据。

使用场景:

  • 适用于写入频繁、单次插入数据量小且对实时性要求不高的场景,例如日志数据插入。

从作用、重要性和影响、使用场景三个方面来看,

  1. insert_distributed_sync 针对分布式表的数据分片同步。async_insert 是针对插入请求的处理方式,不区分表类型(可以是本地表也可以是分布式表)。
  2. insert_distributed_sync 控制数据在集群分片间的同步写入。async_insert 控制是否将多个插入请求合并成批处理。
  3. insert_distributed_sync=1 保证数据同步到所有分片,强一致性。async_insert=1 可能会因为批量处理的延迟而导致数据不能立即可见,且存在故障时丢失数据的风险
  4. insert_distributed_sync=1 会降低分布式表插入的吞吐量,但保证一致性。async_insert=1 会提高插入吞吐量,但牺牲了实时性和一定的可靠性。

二、参数配置

全局用户基本设置

users.xml配置文件里

<insert_distributed_sync>1</insert_distributed_sync>

在这里插入图片描述

会话基本session

hb3 :) set insert_distributed_sync=1;

SET insert_distributed_sync = 1

Query id: 5dbdce47-4f39-454f-b3f7-1940edd1070e

Ok.

0 rows in set. Elapsed: 0.001 sec.

hb3 :) select * from system.settings where name like '%insert_distributed_sync%' \G

SELECT *
FROM system.settings
WHERE name LIKE '%insert_distributed_sync%'

Query id: cdffefb5-d0e3-49d3-96b5-2da065e1332a

Row 1:
──────
name:        insert_distributed_sync
value:       1
changed:     1
description: If setting is enabled, insert query into distributed waits until data will be sent to all nodes in cluster.
min:         ᴺᵁᴸᴸ
max:         ᴺᵁᴸᴸ
readonly:    0
type:        Bool
default:     0
alias_for:

1 row in set. Elapsed: 0.003 sec.

JDBC连接设置insert_distributed_sync

在这里插入图片描述

import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
/**
 * @author duhj
 * @datete 2025/10/20 10:20
 */
public class ClickhouseUrlTest2 {
    public static void main(String[] args) throws Exception{
        Connection conn = DriverManager.getConnection("jdbc:clickhouse://hb3:8123/default?insert_distributed_sync=1"
        , "default", "wobushidashen");
        PreparedStatement ps = conn.prepareStatement("insert into test.log_test values(?,?)");
        for (int i = 0; i < 1000000; i++){

            ps.setInt(1, i);
            ps.setString(2, "xxx"+i);
            ps.execute();
            Thread.sleep(200);
        }
        ps.close();
        conn.close();
    }
}

URL:jdbc:clickhouse://hb3:8123/default?insert_distributed_sync=1
如果还有额外的参数利用&拼接即可。
例如:jdbc:clickhouse://hb3:8123/default?insert_distributed_sync=1&async_insert=0

三、扩展部分

  1. 源码1分布式表写入时consume方法解释
    在这里插入图片描述
    a.Processor 向 Sink 传入一个 Chunk(块),consume 做首块延迟策略检查然后决定同步/异步写入
    b.首次 chunk 时调用 storage.delayInsertOrThrowIfNeeded()(可基于 local queue 大小延迟或抛错)
    c.将 Chunk 转为 Block(使用 header)
    d.如果 insert_sync 为 true 调用 writeSync,否则 writeAsync

  2. 源码2同步写入block部分
    作用:同步插入主逻辑,针对多个 shard/replica 并发发送并等待完成

void DistributedSink::writeSync(const Block & block)
{
    OpenTelemetrySpanHolder span(__PRETTY_FUNCTION__);
    获取基本配置信息
    const Settings & settings = context->getSettingsRef();
    const auto & shards_info = cluster->getShardsInfo();
    Block block_to_send = removeSuperfluousColumns(block);
    size_t start = 0;
    size_t end = shards_info.size();
    根据分片数计算star、end
    if (settings.insert_shard_id)
    {
        start = settings.insert_shard_id - 1;
        end = settings.insert_shard_id;
    }
    创建线程池pool
    if (!pool)
    {
        /// Deferred initialization. Only for sync insertion.
        initWritingJobs(block_to_send, start, end);
        如果分片键是random那么值=1,否则等于远程和本地任务总和
        size_t jobs_count = random_shard_insert ? 1 : (remote_jobs_count + local_jobs_count);
        最大线程数
        size_t max_threads = std::min<size_t>(settings.max_distributed_connections, jobs_count);
        计算最大线程数,取设置中的最大分布式连接数和任务数的最小值
        pool.emplace(/* max_threads_= */ max_threads,
                     /* max_free_threads_= */ max_threads,
                     /* queue_size_= */ jobs_count);

        限流器
        if (!throttler && (settings.max_network_bandwidth || settings.max_network_bytes))
        {
            throttler = std::make_shared<Throttler>(settings.max_network_bandwidth, settings.max_network_bytes,
                                                    "Network bandwidth limit for a query exceeded.");
        }
        重启计时器watch(用于监控整体时间)
        watch.restart();
    }
    重启当前块的计时器
    watch_current_block.restart();
    if (random_shard_insert)
    {
        start = storage.getRandomShardIndex(shards_info);
        end = start + 1;
    }
    size_t num_shards = end - start;
    span.addAttribute("clickhouse.start_shard", start);
    span.addAttribute("clickhouse.end_shard", end);
    span.addAttribute("db.statement", this->query_string);
    如果分片数量大于1,则创建选择器(selector)并为每个分片准备行号
    if (num_shards > 1)
    {
        auto current_selector = createSelector(block);

        /// Prepare row numbers for needed shards
        for (size_t shard_index : collections::range(start, end))
            per_shard_jobs[shard_index].shard_current_block_permutation.resize(0);

        for (size_t i = 0; i < block.rows(); ++i)
            per_shard_jobs[current_selector[i]].shard_current_block_permutation.push_back(i);
    }
    try
    {
        /// Run jobs in parallel for each block and wait them
        finished_jobs_count = 0;
        并行运行每个块的任务并等待
        for (size_t shard_index : collections::range(start, end))
            for (JobReplica & job : per_shard_jobs[shard_index].replicas_jobs)
                
                pool->scheduleOrThrowOnError(runWritingJob(job, block_to_send, num_shards));
    }
    catch (...)
    {
        pool->wait();
        throw;
    }
    try
    {
        waitForJobs();   等待所有分片执行完job
    }
    catch (Exception & exception)
    {
        exception.addMessage(getCurrentStateDescription());
        span.addAttribute(exception);
        throw;
    }
    更新插入的块数和行数
    inserted_blocks += 1;
    inserted_rows += block.rows();
}

a.准备 block_to_send(removeSuperfluousColumns)
b.延迟初始化 pool(线程池)和 throttler(如果尚未创建)。
c.计算 start/end(根据 settings.insert_shard_id 或 random_shard_insert)。
d.如果 num_shards>1:使用 createSelector 生成 selector 并为每个 shard 构造 shard_current_block_permutation。
e.对每个 shard 的每个 replicas_jobs 调用 pool->scheduleOrThrowOnError(runWritingJob(…)) 提交任务。
f.等待所有任务完成(waitForJobs),处理异常并补充状态信息。
g.成功后统计 inserted_blocks/inserted_rows。

  1. 源码3异步写入源码writeAsync
void DistributedSink::writeAsync(const Block & block)
{
    随机分片写入
    if (random_shard_insert)
    {
        writeAsyncImpl(block, storage.getRandomShardIndex(cluster->getShardsInfo()));
        ++inserted_blocks;
    }
    else
    {
        if (storage.getShardingKeyExpr() && (cluster->getShardsInfo().size() > 1))
            return writeSplitAsync(block);

        writeAsyncImpl(block);
        ++inserted_blocks;
    }
}

a.若 random_shard_insert 为 true:随机挑选一个 shard 并调用 writeAsyncImpl(block, shard_id)
b.否则:如果有分片键且集群分片 >1,则调用 writeSplitAsync(先按 selector 切分块);否则直接 writeAsyncImpl(block)
c.统计 inserted_blocks++

  1. 异步写实现源码部分


void DistributedSink::writeAsyncImpl(const Block & block, size_t shard_id)
{
    OpenTelemetrySpanHolder span("DistributedSink::writeAsyncImpl()");

    const auto & shard_info = cluster->getShardsInfo()[shard_id];
    const auto & settings = context->getSettingsRef();
    Block block_to_send = removeSuperfluousColumns(block);

    span.addAttribute("clickhouse.shard_num", shard_info.shard_num);
    span.addAttribute("clickhouse.written_rows", block.rows());

    if (shard_info.hasInternalReplication())
    {
        if (shard_info.isLocal() && settings.prefer_localhost_replica)
            /// Prefer insert into current instance directly
            writeToLocal(block_to_send, shard_info.getLocalNodeCount());
        else
        {
            const auto & path = shard_info.insertPathForInternalReplication(
                settings.prefer_localhost_replica,
                settings.use_compact_format_in_distributed_parts_names);
            if (path.empty())
                throw Exception("Directory name for async inserts is empty", ErrorCodes::LOGICAL_ERROR);
            writeToShard(block_to_send, {path});
        }
    }
    else
    {
        if (shard_info.isLocal() && settings.prefer_localhost_replica)
            writeToLocal(block_to_send, shard_info.getLocalNodeCount());

        std::vector<std::string> dir_names;
        for (const auto & address : cluster->getShardsAddresses()[shard_id])
            if (!address.is_local || !settings.prefer_localhost_replica)
                dir_names.push_back(address.toFullString(settings.use_compact_format_in_distributed_parts_names));

        if (!dir_names.empty())
            writeToShard(block_to_send, dir_names);
    }
}

更多推荐