2013-11-08 17:43:03 +00:00
|
|
|
#pragma once
|
|
|
|
|
2017-04-01 09:19:00 +00:00
|
|
|
#include <DataStreams/copyData.h>
|
|
|
|
#include <DataStreams/IBlockOutputStream.h>
|
|
|
|
#include <DataStreams/OneBlockInputStream.h>
|
|
|
|
#include <DataStreams/MaterializingBlockInputStream.h>
|
|
|
|
#include <Interpreters/InterpreterSelectQuery.h>
|
2017-07-25 21:07:05 +00:00
|
|
|
#include <Storages/StorageMaterializedView.h>
|
2013-11-08 17:43:03 +00:00
|
|
|
|
|
|
|
|
|
|
|
namespace DB
|
|
|
|
{
|
|
|
|
|
|
|
|
|
2017-07-25 21:07:05 +00:00
|
|
|
/** Writes data to the specified table and to all dependent materialized views.
|
2013-11-08 17:43:03 +00:00
|
|
|
*/
|
2013-11-13 14:39:48 +00:00
|
|
|
class PushingToViewsBlockOutputStream : public IBlockOutputStream
|
2013-11-08 17:43:03 +00:00
|
|
|
{
|
|
|
|
public:
|
2017-05-25 01:12:41 +00:00
|
|
|
PushingToViewsBlockOutputStream(String database, String table, const Context & context_, const ASTPtr & query_ptr_)
|
2017-04-01 07:20:54 +00:00
|
|
|
: context(context_), query_ptr(query_ptr_)
|
|
|
|
{
|
|
|
|
storage = context.getTable(database, table);
|
|
|
|
|
2017-05-13 22:19:04 +00:00
|
|
|
/** TODO This is a very important line. At any insertion into the table one of streams should own lock.
|
|
|
|
* Although now any insertion into the table is done via PushingToViewsBlockOutputStream,
|
|
|
|
* but it's clear that here is not the best place for this functionality.
|
2017-04-01 07:20:54 +00:00
|
|
|
*/
|
2017-09-01 15:05:23 +00:00
|
|
|
addTableLock(storage->lockStructure(true, __PRETTY_FUNCTION__));
|
2017-04-01 07:20:54 +00:00
|
|
|
|
|
|
|
Dependencies dependencies = context.getDependencies(database, table);
|
2017-07-25 21:07:05 +00:00
|
|
|
for (const auto & database_table : dependencies)
|
|
|
|
views.emplace_back(
|
|
|
|
dynamic_cast<const StorageMaterializedView &>(*context.getTable(database_table.first, database_table.second)).getInnerQuery(),
|
|
|
|
std::make_shared<PushingToViewsBlockOutputStream>(database_table.first, database_table.second, context, ASTPtr()));
|
2017-04-01 07:20:54 +00:00
|
|
|
|
2017-07-25 21:07:05 +00:00
|
|
|
output = storage->write(query_ptr, context.getSettingsRef());
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
void write(const Block & block) override
|
|
|
|
{
|
2017-07-25 21:07:05 +00:00
|
|
|
for (auto & view : views)
|
2017-04-01 07:20:54 +00:00
|
|
|
{
|
|
|
|
BlockInputStreamPtr from = std::make_shared<OneBlockInputStream>(block);
|
2017-07-25 21:07:05 +00:00
|
|
|
InterpreterSelectQuery select(view.first, context, QueryProcessingStage::Complete, 0, from);
|
2017-04-01 07:20:54 +00:00
|
|
|
BlockInputStreamPtr data = std::make_shared<MaterializingBlockInputStream>(select.execute().in);
|
2017-07-25 21:07:05 +00:00
|
|
|
copyData(*data, *view.second);
|
2017-04-01 07:20:54 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
if (output)
|
|
|
|
output->write(block);
|
|
|
|
}
|
|
|
|
|
|
|
|
void flush() override
|
|
|
|
{
|
|
|
|
if (output)
|
|
|
|
output->flush();
|
|
|
|
}
|
|
|
|
|
|
|
|
void writePrefix() override
|
|
|
|
{
|
|
|
|
if (output)
|
|
|
|
output->writePrefix();
|
|
|
|
}
|
|
|
|
|
|
|
|
void writeSuffix() override
|
|
|
|
{
|
|
|
|
if (output)
|
|
|
|
output->writeSuffix();
|
|
|
|
}
|
2013-11-08 17:43:03 +00:00
|
|
|
|
|
|
|
private:
|
2017-04-01 07:20:54 +00:00
|
|
|
StoragePtr storage;
|
|
|
|
BlockOutputStreamPtr output;
|
2017-09-04 17:49:39 +00:00
|
|
|
const Context & context;
|
2017-04-01 07:20:54 +00:00
|
|
|
ASTPtr query_ptr;
|
2017-07-25 21:07:05 +00:00
|
|
|
std::vector<std::pair<ASTPtr, BlockOutputStreamPtr>> views;
|
2013-11-08 17:43:03 +00:00
|
|
|
};
|
|
|
|
|
|
|
|
|
|
|
|
}
|