StorageReplicatedMergeTree.cpp 19.1 KB
Newer Older
M
Merge  
Michael Kolupaev 已提交
1
#include <DB/Storages/StorageReplicatedMergeTree.h>
M
Merge  
Michael Kolupaev 已提交
2
#include <DB/Storages/MergeTree/ReplicatedMergeTreeBlockOutputStream.h>
M
Merge  
Michael Kolupaev 已提交
3
#include <DB/Storages/MergeTree/ReplicatedMergeTreePartsExchange.h>
M
Merge  
Michael Kolupaev 已提交
4 5 6
#include <DB/Parsers/formatAST.h>
#include <DB/IO/WriteBufferFromOStream.h>
#include <DB/IO/ReadBufferFromString.h>
M
Merge  
Michael Kolupaev 已提交
7 8 9 10

namespace DB
{

M
Merge  
Michael Kolupaev 已提交
11 12 13 14 15 16 17

const auto QUEUE_UPDATE_SLEEP = std::chrono::seconds(5);
const auto QUEUE_NO_WORK_SLEEP = std::chrono::seconds(5);
const auto QUEUE_ERROR_SLEEP = std::chrono::seconds(1);
const auto QUEUE_AFTER_WORK_SLEEP = std::chrono::seconds(0);


M
Merge  
Michael Kolupaev 已提交
18 19 20
StorageReplicatedMergeTree::StorageReplicatedMergeTree(
	const String & zookeeper_path_,
	const String & replica_name_,
M
Merge  
Michael Kolupaev 已提交
21
	bool attach,
M
Merge  
Michael Kolupaev 已提交
22
	const String & path_, const String & name_, NamesAndTypesListPtr columns_,
M
Merge  
Michael Kolupaev 已提交
23
	Context & context_,
M
Merge  
Michael Kolupaev 已提交
24 25 26 27 28 29 30 31
	ASTPtr & primary_expr_ast_,
	const String & date_column_name_,
	const ASTPtr & sampling_expression_,
	size_t index_granularity_,
	MergeTreeData::Mode mode_,
	const String & sign_column_,
	const MergeTreeSettings & settings_)
	:
M
Merge  
Michael Kolupaev 已提交
32
	context(context_), zookeeper(context.getZooKeeper()),
M
Merge  
Michael Kolupaev 已提交
33 34 35 36
	path(path_), name(name_), full_path(path + escapeForFileName(name) + '/'), zookeeper_path(zookeeper_path_),
	replica_name(replica_name_),
	data(	full_path, columns_, context_, primary_expr_ast_, date_column_name_, sampling_expression_,
			index_granularity_,mode_, sign_column_, settings_),
M
Merge  
Michael Kolupaev 已提交
37
	reader(data), writer(data), fetcher(data),
M
Merge  
Michael Kolupaev 已提交
38
	log(&Logger::get("StorageReplicatedMergeTree")),
M
Merge  
Michael Kolupaev 已提交
39 40
	shutdown_called(false)
{
M
Merge  
Michael Kolupaev 已提交
41 42 43
	if (!zookeeper_path.empty() && *zookeeper_path.rbegin() == '/')
		zookeeper_path.erase(zookeeper_path.end() - 1);
	replica_path = zookeeper_path + "/replicas/" + replica_name;
44

M
Merge  
Michael Kolupaev 已提交
45 46
	if (!attach)
	{
M
Merge  
Michael Kolupaev 已提交
47 48 49 50 51 52 53 54
		if (!zookeeper.exists(zookeeper_path))
			createTable();

		if (!isTableEmpty())
			throw Exception("Can't add new replica to non-empty table", ErrorCodes::ADDING_REPLICA_TO_NON_EMPTY_TABLE);

		checkTableStructure();
		createReplica();
M
Merge  
Michael Kolupaev 已提交
55 56 57 58
	}
	else
	{
		checkTableStructure();
M
Merge  
Michael Kolupaev 已提交
59
		checkParts();
M
Merge  
Michael Kolupaev 已提交
60
	}
M
Merge  
Michael Kolupaev 已提交
61

M
Merge  
Michael Kolupaev 已提交
62
	loadQueue();
M
Merge  
Michael Kolupaev 已提交
63
	activateReplica();
M
Merge  
Michael Kolupaev 已提交
64 65 66 67

	queue_updating_thread = std::thread(&StorageReplicatedMergeTree::queueUpdatingThread, this);
	for (size_t i = 0; i < settings_.replication_threads; ++i)
		queue_threads.push_back(std::thread(&StorageReplicatedMergeTree::queueThread, this));
M
Merge  
Michael Kolupaev 已提交
68 69 70 71 72
}

StoragePtr StorageReplicatedMergeTree::create(
	const String & zookeeper_path_,
	const String & replica_name_,
M
Merge  
Michael Kolupaev 已提交
73
	bool attach,
M
Merge  
Michael Kolupaev 已提交
74
	const String & path_, const String & name_, NamesAndTypesListPtr columns_,
M
Merge  
Michael Kolupaev 已提交
75
	Context & context_,
M
Merge  
Michael Kolupaev 已提交
76 77 78 79 80 81 82 83
	ASTPtr & primary_expr_ast_,
	const String & date_column_name_,
	const ASTPtr & sampling_expression_,
	size_t index_granularity_,
	MergeTreeData::Mode mode_,
	const String & sign_column_,
	const MergeTreeSettings & settings_)
{
M
Merge  
Michael Kolupaev 已提交
84 85 86 87
	StorageReplicatedMergeTree * res = new StorageReplicatedMergeTree(zookeeper_path_, replica_name_, attach,
		path_, name_, columns_, context_, primary_expr_ast_, date_column_name_, sampling_expression_,
		index_granularity_, mode_, sign_column_, settings_);
	StoragePtr res_ptr = res->thisPtr();
M
Merge  
Michael Kolupaev 已提交
88 89 90
	String endpoint_name = "ReplicatedMergeTree:" + res->replica_path;
	InterserverIOEndpointPtr endpoint = new ReplicatedMergeTreePartsServer(res->data, res_ptr);
	res->endpoint_holder = new InterserverIOEndpointHolder(endpoint_name, endpoint, res->context.getInterserverIOHandler());
M
Merge  
Michael Kolupaev 已提交
91
	return res_ptr;
M
Merge  
Michael Kolupaev 已提交
92 93
}

M
Merge  
Michael Kolupaev 已提交
94 95 96 97 98 99 100 101
static String formattedAST(const ASTPtr & ast)
{
	if (!ast)
		return "";
	std::stringstream ss;
	formatAST(*ast, ss, 0, false, true);
	return ss.str();
}
M
Merge  
Michael Kolupaev 已提交
102

M
Merge  
Michael Kolupaev 已提交
103
void StorageReplicatedMergeTree::createTable()
M
Merge  
Michael Kolupaev 已提交
104
{
M
Merge  
Michael Kolupaev 已提交
105
	zookeeper.create(zookeeper_path, "", zkutil::CreateMode::Persistent);
M
Merge  
Michael Kolupaev 已提交
106

M
Merge  
Michael Kolupaev 已提交
107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125
	/// Запишем метаданные таблицы, чтобы реплики могли сверять с ними свою локальную структуру таблицы.
	std::stringstream metadata;
	metadata << "metadata format version: 1" << std::endl;
	metadata << "date column: " << data.date_column_name << std::endl;
	metadata << "sampling expression: " << formattedAST(data.sampling_expression) << std::endl;
	metadata << "index granularity: " << data.index_granularity << std::endl;
	metadata << "mode: " << static_cast<int>(data.mode) << std::endl;
	metadata << "sign column: " << data.sign_column << std::endl;
	metadata << "primary key: " << formattedAST(data.primary_expr_ast) << std::endl;
	metadata << "columns:" << std::endl;
	WriteBufferFromOStream buf(metadata);
	for (auto & it : data.getColumnsList())
	{
		writeBackQuotedString(it.first, buf);
		writeChar(' ', buf);
		writeString(it.second->getName(), buf);
		writeChar('\n', buf);
	}
	buf.next();
M
Merge  
Michael Kolupaev 已提交
126

M
Merge  
Michael Kolupaev 已提交
127
	zookeeper.create(zookeeper_path + "/metadata", metadata.str(), zkutil::CreateMode::Persistent);
M
Merge  
Michael Kolupaev 已提交
128

M
Merge  
Michael Kolupaev 已提交
129 130 131
	/// Создадим нужные "директории".
	zookeeper.create(zookeeper_path + "/replicas", "", zkutil::CreateMode::Persistent);
	zookeeper.create(zookeeper_path + "/blocks", "", zkutil::CreateMode::Persistent);
M
Merge  
Michael Kolupaev 已提交
132
	zookeeper.create(zookeeper_path + "/block_numbers", "", zkutil::CreateMode::Persistent);
M
Merge  
Michael Kolupaev 已提交
133
	zookeeper.create(zookeeper_path + "/temp", "", zkutil::CreateMode::Persistent);
M
Merge  
Michael Kolupaev 已提交
134
}
M
Merge  
Michael Kolupaev 已提交
135

M
Merge  
Michael Kolupaev 已提交
136 137
/** Проверить, что список столбцов и настройки таблицы совпадают с указанными в ZK (/metadata).
	* Если нет - бросить исключение.
M
Merge  
Michael Kolupaev 已提交
138
	*/
M
Merge  
Michael Kolupaev 已提交
139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167
void StorageReplicatedMergeTree::checkTableStructure()
{
	String metadata_str = zookeeper.get(zookeeper_path + "/metadata");
	ReadBufferFromString buf(metadata_str);
	assertString("metadata format version: 1", buf);
	assertString("\ndate column: ", buf);
	assertString(data.date_column_name, buf);
	assertString("\nsampling expression: ", buf);
	assertString(formattedAST(data.sampling_expression), buf);
	assertString("\nindex granularity: ", buf);
	assertString(toString(data.index_granularity), buf);
	assertString("\nmode: ", buf);
	assertString(toString(static_cast<int>(data.mode)), buf);
	assertString("\nsign column: ", buf);
	assertString(data.sign_column, buf);
	assertString("\nprimary key: ", buf);
	assertString(formattedAST(data.primary_expr_ast), buf);
	assertString("\ncolumns:\n", buf);
	for (auto & it : data.getColumnsList())
	{
		String name;
		readBackQuotedString(name, buf);
		if (name != it.first)
			throw Exception("Unexpected column name in ZooKeeper: expected " + it.first + ", found " + name,
				ErrorCodes::UNKNOWN_IDENTIFIER);
		assertString(" ", buf);
		assertString(it.second->getName(), buf);
		assertString("\n", buf);
	}
M
Merge  
Michael Kolupaev 已提交
168
	assertEOF(buf);
M
Merge  
Michael Kolupaev 已提交
169
}
M
Merge  
Michael Kolupaev 已提交
170

M
Merge  
Michael Kolupaev 已提交
171 172 173 174 175 176 177 178 179
void StorageReplicatedMergeTree::createReplica()
{
	zookeeper.create(replica_path, "", zkutil::CreateMode::Persistent);
	zookeeper.create(replica_path + "/host", "", zkutil::CreateMode::Persistent);
	zookeeper.create(replica_path + "/log", "", zkutil::CreateMode::Persistent);
	zookeeper.create(replica_path + "/log_pointers", "", zkutil::CreateMode::Persistent);
	zookeeper.create(replica_path + "/queue", "", zkutil::CreateMode::Persistent);
	zookeeper.create(replica_path + "/parts", "", zkutil::CreateMode::Persistent);
}
M
Merge  
Michael Kolupaev 已提交
180

M
Merge  
Michael Kolupaev 已提交
181 182 183 184 185
void StorageReplicatedMergeTree::activateReplica()
{
	std::stringstream host;
	host << "host: " << context.getInterserverIOHost() << std::endl;
	host << "port: " << context.getInterserverIOPort() << std::endl;
M
Merge  
Michael Kolupaev 已提交
186

M
Merge  
Michael Kolupaev 已提交
187
	/// Одновременно объявим, что эта реплика активна, и обновим хост.
M
Merge  
Michael Kolupaev 已提交
188 189 190
	zkutil::Ops ops;
	ops.push_back(new zkutil::Op::Create(replica_path + "/is_active", "", zookeeper.getDefaultACL(), zkutil::CreateMode::Ephemeral));
	ops.push_back(new zkutil::Op::SetData(replica_path + "/host", host.str(), -1));
M
Merge  
Michael Kolupaev 已提交
191 192 193 194 195 196 197 198 199 200 201 202 203

	try
	{
		zookeeper.multi(ops);
	}
	catch (zkutil::KeeperException & e)
	{
		if (e.code == zkutil::ReturnCode::NodeExists)
			throw Exception("Replica " + replica_path + " appears to be already active. If you're sure it's not, "
				"try again in a minute or remove znode " + replica_path + "/is_active manually", ErrorCodes::REPLICA_IS_ALREADY_ACTIVE);

		throw;
	}
M
Merge  
Michael Kolupaev 已提交
204

M
Merge  
Michael Kolupaev 已提交
205 206
	replica_is_active_node = zkutil::EphemeralNodeHolder::existing(replica_path + "/is_active", zookeeper);
}
M
Merge  
Michael Kolupaev 已提交
207

M
Merge  
Michael Kolupaev 已提交
208 209 210 211 212 213 214 215 216 217
bool StorageReplicatedMergeTree::isTableEmpty()
{
	Strings replicas = zookeeper.getChildren(zookeeper_path + "/replicas");
	for (const auto & replica : replicas)
	{
		if (!zookeeper.getChildren(zookeeper_path + "/replicas/" + replica + "/parts").empty())
			return false;
	}
	return true;
}
M
Merge  
Michael Kolupaev 已提交
218

M
Merge  
Michael Kolupaev 已提交
219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240
void StorageReplicatedMergeTree::checkParts()
{
	Strings expected_parts_vec = zookeeper.getChildren(replica_path + "/parts");
	NameSet expected_parts(expected_parts_vec.begin(), expected_parts_vec.end());

	MergeTreeData::DataParts parts = data.getDataParts();

	MergeTreeData::DataPartsVector unexpected_parts;
	for (const auto & part : parts)
	{
		if (expected_parts.count(part->name))
		{
			expected_parts.erase(part->name);
		}
		else
		{
			unexpected_parts.push_back(part);
		}
	}

	if (!expected_parts.empty())
		throw Exception("Not found " + toString(expected_parts.size())
M
Merge  
Michael Kolupaev 已提交
241
			+ " parts (including " + *expected_parts.begin() + ") in table " + data.getTableName(),
M
Merge  
Michael Kolupaev 已提交
242 243 244 245 246 247 248 249 250 251
			ErrorCodes::NOT_FOUND_EXPECTED_DATA_PART);

	if (unexpected_parts.size() > 1)
		throw Exception("More than one unexpected part (including " + unexpected_parts[0]->name
			+ ") in table " + data.getTableName(),
			ErrorCodes::TOO_MANY_UNEXPECTED_DATA_PARTS);

	for (MergeTreeData::DataPartPtr part : unexpected_parts)
	{
		LOG_ERROR(log, "Unexpected part " << part->name << ". Renaming it to ignored_" + part->name);
M
Merge  
Michael Kolupaev 已提交
252
		data.renameAndDetachPart(part, "ignored_");
M
Merge  
Michael Kolupaev 已提交
253 254
	}
}
M
Merge  
Michael Kolupaev 已提交
255

M
Merge  
Michael Kolupaev 已提交
256 257 258 259 260 261 262 263
void StorageReplicatedMergeTree::loadQueue()
{
	Poco::ScopedLock<Poco::FastMutex> lock(queue_mutex);

	Strings children = zookeeper.getChildren(replica_path + "/queue");
	std::sort(children.begin(), children.end());
	for (const String & child : children)
	{
M
Merge  
Michael Kolupaev 已提交
264
		String s = zookeeper.get(replica_path + "/queue/" + child);
M
Merge  
Michael Kolupaev 已提交
265 266 267
		LogEntry entry = LogEntry::parse(s);
		entry.znode_name = child;
		queue.push_back(entry);
M
Merge  
Michael Kolupaev 已提交
268 269 270 271 272 273
	}
}

void StorageReplicatedMergeTree::pullLogsToQueue()
{
	Poco::ScopedLock<Poco::FastMutex> lock(queue_mutex);
M
Merge  
Michael Kolupaev 已提交
274

M
Merge  
Michael Kolupaev 已提交
275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305
	/// Сольем все логи в хронологическом порядке.

	struct LogIterator
	{
		String replica;		/// Имя реплики.
		UInt64 index;		/// Номер записи в логе (суффикс имени ноды).

		Int64 timestamp;	/// Время (czxid) создания записи в логе.
		String entry_str;	/// Сама запись.

		bool operator<(const LogIterator & rhs) const
		{
			return timestamp < rhs.timestamp;
		}

		bool readEntry(zkutil::ZooKeeper & zookeeper, const String & zookeeper_path)
		{
			String index_str = toString(index);
			while (index_str.size() < 10)
				index_str = '0' + index_str;
			zkutil::Stat stat;
			if (!zookeeper.tryGet(zookeeper_path + "/replicas/" + replica + "/log/log-" + index_str, entry_str, &stat))
				return false;
			timestamp = stat.getczxid();
			return true;
		}
	};

	typedef std::priority_queue<LogIterator> PriorityQueue;
	PriorityQueue priority_queue;

M
Merge  
Michael Kolupaev 已提交
306
	Strings replicas = zookeeper.getChildren(zookeeper_path + "/replicas");
M
Merge  
Michael Kolupaev 已提交
307

M
Merge  
Michael Kolupaev 已提交
308 309
	for (const String & replica : replicas)
	{
M
Merge  
Michael Kolupaev 已提交
310 311
		String index_str;
		UInt64 index;
M
Merge  
Michael Kolupaev 已提交
312

M
Merge  
Michael Kolupaev 已提交
313
		if (zookeeper.tryGet(replica_path + "/log_pointers/" + replica, index_str))
M
Merge  
Michael Kolupaev 已提交
314
		{
M
Merge  
Michael Kolupaev 已提交
315
			index = Poco::NumberParser::parseUnsigned64(index_str);
M
Merge  
Michael Kolupaev 已提交
316 317 318 319
		}
		else
		{
			/// Если у нас еще нет указателя на лог этой реплики, поставим указатель на первую запись в нем.
M
Merge  
Michael Kolupaev 已提交
320
			Strings entries = zookeeper.getChildren(zookeeper_path + "/replicas/" + replica + "/log");
M
Merge  
Michael Kolupaev 已提交
321
			std::sort(entries.begin(), entries.end());
M
Merge  
Michael Kolupaev 已提交
322
			index = entries.empty() ? 0 : Poco::NumberParser::parseUnsigned64(entries[0].substr(strlen("log-")));
M
Merge  
Michael Kolupaev 已提交
323

M
Merge  
Michael Kolupaev 已提交
324
			zookeeper.create(replica_path + "/log_pointers/" + replica, toString(index), zkutil::CreateMode::Persistent);
M
Merge  
Michael Kolupaev 已提交
325 326
		}

M
Merge  
Michael Kolupaev 已提交
327 328 329 330 331 332 333
		LogIterator iterator;
		iterator.replica = replica;
		iterator.index = index;

		if (iterator.readEntry(zookeeper, zookeeper_path))
			priority_queue.push(iterator);
	}
M
Merge  
Michael Kolupaev 已提交
334

M
Merge  
Michael Kolupaev 已提交
335 336 337 338
	while (!priority_queue.empty())
	{
		LogIterator iterator = priority_queue.top();
		priority_queue.pop();
M
Merge  
Michael Kolupaev 已提交
339

M
Merge  
Michael Kolupaev 已提交
340
		LogEntry entry = LogEntry::parse(iterator.entry_str);
M
Merge  
Michael Kolupaev 已提交
341

M
Merge  
Michael Kolupaev 已提交
342 343 344 345 346 347 348 349 350 351 352 353 354 355 356
		/// Одновременно добавим запись в очередь и продвинем указатель на лог.
		zkutil::Ops ops;
		ops.push_back(new zkutil::Op::Create(
			replica_path + "/queue/queue-", iterator.entry_str, zookeeper.getDefaultACL(), zkutil::CreateMode::PersistentSequential));
		ops.push_back(new zkutil::Op::SetData(
			replica_path + "/log_pointers/" + iterator.replica, toString(iterator.index + 1), -1));
		auto results = zookeeper.multi(ops);

		String path_created = dynamic_cast<zkutil::OpResult::Create &>((*results)[0]).getPathCreated();
		entry.znode_name = path.substr(path.find_last_of('/') + 1);
		queue.push_back(entry);

		++iterator.index;
		if (iterator.readEntry(zookeeper, zookeeper_path))
			priority_queue.push(iterator);
M
Merge  
Michael Kolupaev 已提交
357 358 359 360 361
	}
}

void StorageReplicatedMergeTree::optimizeQueue()
{
M
Merge  
Michael Kolupaev 已提交
362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378
}

void StorageReplicatedMergeTree::executeLogEntry(const LogEntry & entry)
{
	if (entry.type == LogEntry::GET_PART ||
		entry.type == LogEntry::MERGE_PARTS)
	{
		/// Если у нас уже есть этот кусок или покрывающий его кусок, ничего делать не нужно.
		MergeTreeData::DataPartPtr containing_part = data.getContainingPart(entry.new_part_name);

		/// Даже если кусок есть локально, его (в исключительных случаях) может не быть в zookeeper.
		if (containing_part && zookeeper.exists(replica_path + "/parts/" + containing_part->name))
			return;
	}

	if (entry.type != LogEntry::GET_PART)
		throw Exception("Merging is not implemented.", ErrorCodes::NOT_IMPLEMENTED);
M
Merge  
Michael Kolupaev 已提交
379

M
Merge  
Michael Kolupaev 已提交
380 381 382 383 384 385 386 387 388 389 390 391
	String replica = findActiveReplicaHavingPart(entry.new_part_name);
	fetchPart(entry.new_part_name, replica);
}

void StorageReplicatedMergeTree::queueUpdatingThread()
{
	while (!shutdown_called)
	{
		pullLogsToQueue();
		optimizeQueue();
		std::this_thread::sleep_for(QUEUE_UPDATE_SLEEP);
	}
M
Merge  
Michael Kolupaev 已提交
392
}
M
Merge  
Michael Kolupaev 已提交
393

M
Merge  
Michael Kolupaev 已提交
394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409
void StorageReplicatedMergeTree::queueThread()
{
	while (!shutdown_called)
	{
		LogEntry entry;
		bool empty;

		{
			Poco::ScopedLock<Poco::FastMutex> lock(queue_mutex);
			empty = queue.empty();
			if (!empty)
			{
				entry = queue.front();
				queue.pop_front();
			}
		}
M
Merge  
Michael Kolupaev 已提交
410

M
Merge  
Michael Kolupaev 已提交
411 412 413 414 415
		if (empty)
		{
			std::this_thread::sleep_for(QUEUE_NO_WORK_SLEEP);
			continue;
		}
M
Merge  
Michael Kolupaev 已提交
416

M
Merge  
Michael Kolupaev 已提交
417
		bool success = false;
M
Merge  
Michael Kolupaev 已提交
418

M
Merge  
Michael Kolupaev 已提交
419 420 421
		try
		{
			executeLogEntry(entry);
M
Merge  
Michael Kolupaev 已提交
422
			zookeeper.remove(replica_path + "/queue/" + entry.znode_name);
M
Merge  
Michael Kolupaev 已提交
423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497

			success = true;
		}
		catch (const Exception & e)
		{
			LOG_ERROR(log, "Code: " << e.code() << ". " << e.displayText() << std::endl
				<< std::endl
				<< "Stack trace:" << std::endl
				<< e.getStackTrace().toString());
		}
		catch (const Poco::Exception & e)
		{
			LOG_ERROR(log, "Poco::Exception: " << e.code() << ". " << e.displayText());
		}
		catch (const std::exception & e)
		{
			LOG_ERROR(log, "std::exception: " << e.what());
		}
		catch (...)
		{
			LOG_ERROR(log, "Unknown exception");
		}

		if (shutdown_called)
			break;

		if (success)
		{
			std::this_thread::sleep_for(QUEUE_AFTER_WORK_SLEEP);
		}
		else
		{
			{
				/// Добавим действие, которое не получилось выполнить, в конец очереди.
				Poco::ScopedLock<Poco::FastMutex> lock(queue_mutex);
				queue.push_back(entry);
			}
			std::this_thread::sleep_for(QUEUE_ERROR_SLEEP);
		}
	}
}

String StorageReplicatedMergeTree::findActiveReplicaHavingPart(const String & part_name)
{
	Strings replicas = zookeeper.getChildren(zookeeper_path + "/replicas");

	/// Из реплик, у которых есть кусок, выберем одну равновероятно.
	std::random_shuffle(replicas.begin(), replicas.end());

	for (const String & replica : replicas)
	{
		if (zookeeper.exists(zookeeper_path + "/replicas/" + replica + "/parts/" + part_name) &&
			zookeeper.exists(zookeeper_path + "/replicas/" + replica + "/is_active"))
			return replica;
	}

	throw Exception("No active replica has part " + part_name, ErrorCodes::NO_REPLICA_HAS_PART);
}

void StorageReplicatedMergeTree::fetchPart(const String & part_name, const String & replica_name)
{
	auto table_lock = lockStructure(true);

	String host;
	int port;

	String host_port_str = zookeeper.get(zookeeper_path + "/replicas/" + replica_name + "/host");
	ReadBufferFromString buf(host_port_str);
	assertString("host: ", buf);
	readString(host, buf);
	assertString("\nport: ", buf);
	readText(port, buf);
	assertString("\n", buf);
	assertEOF(buf);

M
Merge  
Michael Kolupaev 已提交
498
	MergeTreeData::MutableDataPartPtr part = fetcher.fetchPart(part_name, zookeeper_path + "/replicas/" + replica_name, host, port);
M
Merge  
Michael Kolupaev 已提交
499 500 501 502 503 504 505 506 507 508 509 510 511 512 513
	data.renameTempPartAndAdd(part, nullptr);

	zkutil::Ops ops;
	ops.push_back(new zkutil::Op::Create(
		replica_path + "/parts/" + part->name,
		"",
		zookeeper.getDefaultACL(),
		zkutil::CreateMode::Persistent));
	ops.push_back(new zkutil::Op::Create(
		replica_path + "/parts/" + part->name + "/checksums",
		part->checksums.toString(),
		zookeeper.getDefaultACL(),
		zkutil::CreateMode::Persistent));
	zookeeper.multi(ops);
}
M
Merge  
Michael Kolupaev 已提交
514

M
Merge  
Michael Kolupaev 已提交
515 516 517 518 519 520 521 522
void StorageReplicatedMergeTree::shutdown()
{
	if (shutdown_called)
		return;
	shutdown_called = true;
	replica_is_active_node = nullptr;
	endpoint_holder = nullptr;

M
Merge  
Michael Kolupaev 已提交
523 524 525 526 527
	LOG_TRACE(log, "Waiting for threads to finish");
	queue_updating_thread.join();
	for (auto & thread : queue_threads)
		thread.join();
	LOG_TRACE(log, "Threads finished");
M
Merge  
Michael Kolupaev 已提交
528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552
}

StorageReplicatedMergeTree::~StorageReplicatedMergeTree()
{
	try
	{
		shutdown();
	}
	catch(...)
	{
		tryLogCurrentException("~StorageReplicatedMergeTree");
	}
}

BlockInputStreams StorageReplicatedMergeTree::read(
		const Names & column_names,
		ASTPtr query,
		const Settings & settings,
		QueryProcessingStage::Enum & processed_stage,
		size_t max_block_size,
		unsigned threads)
{
	return reader.read(column_names, query, settings, processed_stage, max_block_size, threads);
}

M
Merge  
Michael Kolupaev 已提交
553 554 555 556 557 558 559 560
BlockOutputStreamPtr StorageReplicatedMergeTree::write(ASTPtr query)
{
	String insert_id;
	if (ASTInsertQuery * insert = dynamic_cast<ASTInsertQuery *>(&*query))
		insert_id = insert->insert_id;

	return new ReplicatedMergeTreeBlockOutputStream(*this, insert_id);
}
M
Merge  
Michael Kolupaev 已提交
561 562 563

void StorageReplicatedMergeTree::drop()
{
M
Merge  
Michael Kolupaev 已提交
564 565
	shutdown();

M
Merge  
Michael Kolupaev 已提交
566 567 568 569 570 571
	replica_is_active_node = nullptr;
	zookeeper.removeRecursive(replica_path);
	if (zookeeper.getChildren(zookeeper_path + "/replicas").empty())
		zookeeper.removeRecursive(zookeeper_path);
}

M
Merge  
Michael Kolupaev 已提交
572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624
void StorageReplicatedMergeTree::LogEntry::writeText(WriteBuffer & out) const
{
	writeString("format version: 1\n", out);
	switch (type)
	{
		case GET_PART:
			writeString("get\n", out);
			writeString(new_part_name, out);
			break;
		case MERGE_PARTS:
			writeString("merge\n", out);
			for (const String & s : parts_to_merge)
			{
				writeString(s, out);
				writeString("\n", out);
			}
			writeString("into\n", out);
			writeString(new_part_name, out);
			break;
	}
	writeString("\n", out);
}

void StorageReplicatedMergeTree::LogEntry::readText(ReadBuffer & in)
{
	String type_str;

	assertString("format version: 1\n", in);
	readString(type_str, in);
	assertString("\n", in);

	if (type_str == "get")
	{
		type = GET_PART;
		readString(new_part_name, in);
	}
	else if (type_str == "merge")
	{
		type = MERGE_PARTS;
		while (true)
		{
			String s;
			readString(s, in);
			assertString("\n", in);
			if (s == "into")
				break;
			parts_to_merge.push_back(s);
		}
		readString(new_part_name, in);
	}
	assertString("\n", in);
}

M
Merge  
Michael Kolupaev 已提交
625
}