* Copyright (c) Huawei Technologies Co., Ltd. 2025. All rights reserved.
* You can use this software according to the terms and conditions of the Mulan PSL v2.
* You may obtain a copy of Mulan PSL v2 at:
* http://license.coscl.org.cn/MulanPSL2
* THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND,
* EITHER EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT,
* MERCHANTABILITY OR FIT FOR A PARTICULAR PURPOSE.
* See the Mulan PSL v2 for more details.
*/
#include "LookupJoinRunner.h"
void LookupJoinRunner::processBatch(omnistream::VectorBatch* in, Context* cxt, Collector* out)
{
collector->setCollector(out);
collector->setInput(in);
collector->reset();
fetcher->flatMap(in, collector);
if (isLeftOuterJoin && !collector->isCollected()) {
NOT_IMPL_EXCEPTION;
}
}
void LookupJoinRunner::open(const Configuration& parameters)
{
using namespace omniruntime::type;
std::string temporalTableSourceSpec = description["temporalTableSourceSpec"];
std::string connectorType = description["connectorType"];
std::string filepath;
if (connectorType == "filesystem") {
filepath = description["connectorPath"];
} else {
NOT_IMPL_EXCEPTION;
}
auto* src = new CsvTableSource(filepath, description["lookupInputTypes"].get<std::vector<std::string>>());
collector = new TableFunctionCollector();
collector->setCollector(innerCollector);
auto lookupFunction = new CsvLookupFunction<int64_t>(description, src);
fetcher = new GeneratedCsvLookupFunction<int64_t>(lookupFunction);
fetcher->open();
}
LookupJoinRunner::LookupJoinRunner(nlohmann::json description, Collector* innerCollector)
: innerCollector(innerCollector),
description(description)
{
isLeftOuterJoin = (description["joinType"].get<std::string>() == "LeftOuterJoin");
}
void LookupJoinRunner::close()
{
delete fetcher;
}