ClickHouse/dbms/include/DB/DataStreams/PushingToViewsBlockOutputStream.h

86 lines
2.7 KiB
C++
Raw Normal View History

#pragma once
#include <DB/DataStreams/copyData.h>
#include <DB/DataStreams/IBlockOutputStream.h>
#include <DB/DataStreams/OneBlockInputStream.h>
#include <DB/DataStreams/MaterializingBlockInputStream.h>
#include <DB/Interpreters/InterpreterSelectQuery.h>
#include <DB/Storages/StorageView.h>
namespace DB
{
/** Записывает данные в указанную таблицу, при этом рекурсивно вызываясь от всех зависимых вьюшек.
* Если вьюшка не материализованная, то в нее данные не записываются, лишь перенаправляются дальше.
*/
class PushingToViewsBlockOutputStream : public IBlockOutputStream
{
public:
PushingToViewsBlockOutputStream(String database, String table, const Context & context_, ASTPtr query_ptr_)
: context(context_), query_ptr(query_ptr_)
{
storage = context.getTable(database, table);
2015-09-18 00:46:36 +00:00
/** TODO Это очень важная строчка. При любой вставке в таблицу один из stream-ов должен владеть lock-ом.
* Хотя сейчас любая вставка в таблицу делается через PushingToViewsBlockOutputStream,
* но ясно, что здесь - не лучшее место для этой функциональности.
*/
addTableLock(storage->lockStructure(true));
Dependencies dependencies = context.getDependencies(database, table);
for (size_t i = 0; i < dependencies.size(); ++i)
{
children.push_back(new PushingToViewsBlockOutputStream(dependencies[i].first, dependencies[i].second, context, ASTPtr()));
queries.push_back(dynamic_cast<StorageView &>(*context.getTable(dependencies[i].first, dependencies[i].second)).getInnerQuery());
}
if (storage->getName() != "View")
output = storage->write(query_ptr, context.getSettingsRef());
}
void write(const Block & block) override
{
for (size_t i = 0; i < children.size(); ++i)
{
BlockInputStreamPtr from = new OneBlockInputStream(block);
InterpreterSelectQuery select(queries[i], context, QueryProcessingStage::Complete, 0, from);
2015-06-18 02:11:05 +00:00
BlockInputStreamPtr data = new MaterializingBlockInputStream(select.execute().in);
copyData(*data, *children[i]);
}
if (output)
output->write(block);
}
2015-01-27 00:52:03 +00:00
void flush() override
{
if (output)
output->flush();
}
void writePrefix() override
{
if (output)
output->writePrefix();
}
void writeSuffix() override
{
if (output)
output->writeSuffix();
}
private:
StoragePtr storage;
BlockOutputStreamPtr output;
Context context;
ASTPtr query_ptr;
std::vector<BlockOutputStreamPtr> children;
std::vector<ASTPtr> queries;
};
}