| 45 | using namespace bytedance::bolt; |
| 46 | namespace bytedance::bolt::tool::trace { |
| 47 | OperatorReplayerBase::OperatorReplayerBase( |
| 48 | std::string traceDir, |
| 49 | std::string queryId, |
| 50 | std::string taskId, |
| 51 | std::string nodeId, |
| 52 | std::string operatorType, |
| 53 | const std::string& driverIds, |
| 54 | uint64_t queryCapacity, |
| 55 | folly::Executor* executor) |
| 56 | : queryId_(std::string(std::move(queryId))), |
| 57 | taskId_(std::move(taskId)), |
| 58 | nodeId_(std::move(nodeId)), |
| 59 | operatorType_(std::move(operatorType)), |
| 60 | taskTraceDir_( |
| 61 | exec::trace::getTaskTraceDirectory(traceDir, queryId_, taskId_)), |
| 62 | nodeTraceDir_(exec::trace::getNodeTraceDirectory(taskTraceDir_, nodeId_)), |
| 63 | fs_(filesystems::getFileSystem(taskTraceDir_, nullptr)), |
| 64 | pipelineIds_(exec::trace::listPipelineIds(nodeTraceDir_, fs_)), |
| 65 | driverIds_( |
| 66 | driverIds.empty() ? exec::trace::listDriverIds( |
| 67 | nodeTraceDir_, |
| 68 | pipelineIds_.front(), |
| 69 | fs_) |
| 70 | : exec::trace::extractDriverIds(driverIds)), |
| 71 | queryCapacity_(queryCapacity == 0 ? memory::kMaxMemory : queryCapacity), |
| 72 | executor_(executor) { |
| 73 | BOLT_USER_CHECK(!taskTraceDir_.empty()); |
| 74 | BOLT_USER_CHECK(!taskId_.empty()); |
| 75 | BOLT_USER_CHECK(!nodeId_.empty()); |
| 76 | BOLT_USER_CHECK(!operatorType_.empty()); |
| 77 | if (operatorType_ == "HashJoin") { |
| 78 | BOLT_USER_CHECK_EQ(pipelineIds_.size(), 2); |
| 79 | } else { |
| 80 | BOLT_USER_CHECK_EQ(pipelineIds_.size(), 1); |
| 81 | } |
| 82 | BOLT_CHECK_NOT_NULL(executor_); |
| 83 | |
| 84 | const auto taskMetaReader = exec::trace::TaskTraceMetadataReader( |
| 85 | taskTraceDir_, memory::MemoryManager::getInstance()->tracePool()); |
| 86 | queryConfigs_ = taskMetaReader.queryConfigs(); |
| 87 | connectorConfigs_ = taskMetaReader.connectorProperties(); |
| 88 | planFragment_ = taskMetaReader.queryPlan(); |
| 89 | queryConfigs_[core::QueryConfig::kQueryTraceEnabled] = "false"; |
| 90 | } |
| 91 | |
| 92 | RowVectorPtr OperatorReplayerBase::run(bool copyResults) { |
| 93 | auto queryCtx = createQueryCtx(); |
nothing calls this directly
no test coverage detected