mirror of
https://github.com/ClickHouse/ClickHouse.git
synced 2024-12-15 02:41:59 +00:00
49 lines
1.1 KiB
C++
49 lines
1.1 KiB
C++
#pragma once
|
|
|
|
#include <Processors/Sinks/SinkToStorage.h>
|
|
#include <Storages/StorageInMemoryMetadata.h>
|
|
|
|
|
|
namespace DB
|
|
{
|
|
|
|
class Block;
|
|
class StorageMergeTree;
|
|
struct StorageSnapshot;
|
|
using StorageSnapshotPtr = std::shared_ptr<StorageSnapshot>;
|
|
|
|
|
|
class MergeTreeSink : public SinkToStorage
|
|
{
|
|
public:
|
|
MergeTreeSink(
|
|
StorageMergeTree & storage_,
|
|
StorageMetadataPtr metadata_snapshot_,
|
|
size_t max_parts_per_block_,
|
|
ContextPtr context_);
|
|
|
|
~MergeTreeSink() override;
|
|
|
|
String getName() const override { return "MergeTreeSink"; }
|
|
void consume(Chunk chunk) override;
|
|
void onStart() override;
|
|
void onFinish() override;
|
|
|
|
private:
|
|
StorageMergeTree & storage;
|
|
StorageMetadataPtr metadata_snapshot;
|
|
size_t max_parts_per_block;
|
|
ContextPtr context;
|
|
StorageSnapshotPtr storage_snapshot;
|
|
UInt64 chunk_dedup_seqnum = 0; /// input chunk ordinal number in case of dedup token
|
|
UInt64 num_blocks_processed = 0;
|
|
|
|
/// We can delay processing for previous chunk and start writing a new one.
|
|
struct DelayedChunk;
|
|
std::unique_ptr<DelayedChunk> delayed_chunk;
|
|
|
|
void finishDelayedChunk();
|
|
};
|
|
|
|
}
|