diff --git a/be/src/runtime/result_block_buffer.cpp b/be/src/runtime/result_block_buffer.cpp index 11c8c981e81ee3..65fe29b8cc3293 100644 --- a/be/src/runtime/result_block_buffer.cpp +++ b/be/src/runtime/result_block_buffer.cpp @@ -157,6 +157,21 @@ Status ResultBlockBuffer::get_batch(std::shared_ptr) { + size_t batch_bytes = 0; + int64_t batch_rows = 0; + if constexpr (std::is_same_v) { + batch_rows = result->rows(); + batch_bytes = result->bytes(); + } + LOG(WARNING) << fmt::format( + "ArrowFlight ResultBlockBuffer get_batch data, fragment_id={}, packet_seq={}, " + "queue_size_after_pop={}, rows={}, bytes={}, is_close={}, waiting_rpc={}, " + "returned_rows={}, dependency_cnt={}", + print_id(_fragment_id), _packet_num, _result_batch_queue.size(), batch_rows, + batch_bytes, _is_close, _waiting_rpc.size(), _returned_rows.load(), + _result_sink_dependencies.size()); + } RETURN_IF_ERROR(ctx->on_data(result, _packet_num, this)); _packet_num++; return Status::OK(); @@ -166,6 +181,14 @@ Status ResultBlockBuffer::get_batch(std::shared_ptron_failure(_status); return Status::OK(); } + if constexpr (std::is_same_v) { + LOG(WARNING) << fmt::format( + "ArrowFlight ResultBlockBuffer get_batch close, fragment_id={}, packet_seq={}, " + "queue_size={}, is_close={}, returned_rows={}, waiting_rpc={}, " + "dependency_cnt={}", + print_id(_fragment_id), _packet_num, _result_batch_queue.size(), _is_close, + _returned_rows.load(), _waiting_rpc.size(), _result_sink_dependencies.size()); + } ctx->on_close(_packet_num, _returned_rows); LOG(INFO) << fmt::format( "ResultBlockBuffer finished, fragment_id={}, is_close={}, is_cancelled={}, " @@ -239,9 +262,33 @@ Status ResultBlockBuffer::add_batch(RuntimeState* state, } _instance_rows[state->fragment_instance_id()] += num_rows; _instance_rows_in_queue.back()[state->fragment_instance_id()] += num_rows; + if constexpr (std::is_same_v) { + LOG(WARNING) << fmt::format( + "ArrowFlight ResultBlockBuffer add_batch queued, fragment_id={}, " + "producer_id={}, " + "rows={}, bytes={}, queue_size={}, packet_seq_next={}, waiting_rpc={}, " + "is_close={}", + print_id(_fragment_id), print_id(state->fragment_instance_id()), num_rows, + batch_size, _result_batch_queue.size(), _packet_num, _waiting_rpc.size(), + _is_close); + } } else { auto ctx = _waiting_rpc.front(); _waiting_rpc.pop_front(); + if constexpr (std::is_same_v) { + int64_t batch_rows = 0; + size_t batch_bytes = 0; + if constexpr (std::is_same_v) { + batch_rows = result->rows(); + batch_bytes = result->bytes(); + } + LOG(WARNING) << fmt::format( + "ArrowFlight ResultBlockBuffer add_batch direct_send, fragment_id={}, " + "producer_id={}, " + "rows={}, bytes={}, packet_seq={}, waiting_rpc_after_pop={}, is_close={}", + print_id(_fragment_id), print_id(state->fragment_instance_id()), batch_rows, + batch_bytes, _packet_num, _waiting_rpc.size(), _is_close); + } RETURN_IF_ERROR(ctx->on_data(result, _packet_num, this)); _packet_num++; } diff --git a/be/src/runtime/result_buffer_mgr.cpp b/be/src/runtime/result_buffer_mgr.cpp index d1c021af9c1a44..612dc0d2509231 100644 --- a/be/src/runtime/result_buffer_mgr.cpp +++ b/be/src/runtime/result_buffer_mgr.cpp @@ -101,6 +101,14 @@ Status ResultBufferMgr::create_sender(const TUniqueId& unique_id, int buffer_siz // add extra 5s for avoid corner case int64_t max_timeout = time(nullptr) + state->execution_timeout() + 5; cancel_at_time(max_timeout, unique_id); + if (arrow_flight) { + LOG(WARNING) << fmt::format( + "ResultBufferMgr create_sender, finstId={}, arrow_flight={}, " + "exec_timeout_s={}, " + "cancel_at={}, buffer_ptr={}, map_size={}", + print_id(unique_id), arrow_flight, state->execution_timeout(), max_timeout, + fmt::ptr(control_block.get()), _buffer_map.size()); + } } *sender = control_block; return Status::OK(); @@ -123,6 +131,13 @@ template Status ResultBufferMgr::find_buffer(const TUniqueId& finst_id, std::shared_ptr& buffer) { buffer = _find_control_block(finst_id); + if constexpr (std::is_same_v) { + const void* buffer_ptr = + buffer != nullptr ? static_cast(buffer.get()) : nullptr; + LOG(WARNING) << fmt::format( + "ResultBufferMgr find_buffer, finstId={}, found={}, buffer_ptr={}", + print_id(finst_id), buffer != nullptr, buffer_ptr); + } return buffer == nullptr ? Status::InternalError( "no arrow schema for this query, maybe query has been " "canceled, finst_id={}", @@ -136,8 +151,22 @@ bool ResultBufferMgr::cancel(const TUniqueId& unique_id, const Status& reason) { auto exist = _buffer_map.end() != iter; if (exist) { + bool arrow_flight = std::dynamic_pointer_cast( + iter->second) != nullptr; + if (arrow_flight) { + LOG(WARNING) << fmt::format( + "ResultBufferMgr cancel, finstId={}, reason={}, buffer_ptr={}, " + "map_size_before={}", + print_id(unique_id), reason.to_string(), fmt::ptr(iter->second.get()), + _buffer_map.size()); + } iter->second->cancel(reason); _buffer_map.erase(iter); + if (arrow_flight) { + LOG(WARNING) << fmt::format( + "ResultBufferMgr cancel erased, finstId={}, map_size_after={}", + print_id(unique_id), _buffer_map.size()); + } } return exist; } diff --git a/be/src/service/arrow_flight/arrow_flight_batch_reader.cpp b/be/src/service/arrow_flight/arrow_flight_batch_reader.cpp index 07e46cfcfed5c3..901688037584ea 100644 --- a/be/src/service/arrow_flight/arrow_flight_batch_reader.cpp +++ b/be/src/service/arrow_flight/arrow_flight_batch_reader.cpp @@ -207,6 +207,12 @@ arrow::Status ArrowFlightBatchRemoteReader::_fetch_data() { while (true) { // if `continue` occurs, data is invalid, continue fetch, block is nullptr. // if `break` occurs, fetch data successfully (block is not nullptr) or fetch eos. + LOG(WARNING) << fmt::format( + "ArrowFlightBatchRemoteReader fetch request, finistId={}, expect_packet_seq={}, " + "result_addr={}:{}, schema_ready={}, cached_block={}, mem_peak={}", + print_id(_statement->query_id), _packet_seq.load(), + _statement->result_addr.hostname, _statement->result_addr.port, _schema != nullptr, + _block != nullptr, _mem_tracker->peak_consumption()); Status st; auto request = std::make_shared(); auto* pfinst_id = request->mutable_finst_id(); @@ -239,6 +245,20 @@ arrow::Status ArrowFlightBatchRemoteReader::_fetch_data() { st = Status::create(callback->response_->status()); ARROW_RETURN_NOT_OK(to_arrow_status(st)); + LOG(WARNING) << fmt::format( + "ArrowFlightBatchRemoteReader fetch response, finistId={}, expect_packet_seq={}, " + "receive_packet_seq={}, has_eos={}, eos={}, has_empty_batch={}, empty_batch={}, " + "has_block={}, block_bytes={}, status={}, result_addr={}:{}", + print_id(_statement->query_id), _packet_seq.load(), + callback->response_->has_packet_seq() ? callback->response_->packet_seq() : -1, + callback->response_->has_eos(), + callback->response_->has_eos() && callback->response_->eos(), + callback->response_->has_empty_batch(), + callback->response_->has_empty_batch() && callback->response_->empty_batch(), + callback->response_->has_block(), + callback->response_->has_block() ? callback->response_->block().ByteSizeLong() : 0, + st.to_string(), _statement->result_addr.hostname, _statement->result_addr.port); + DCHECK(callback->response_->has_packet_seq()); if (_packet_seq != callback->response_->packet_seq()) { return _return_invalid_status( @@ -270,6 +290,11 @@ arrow::Status ArrowFlightBatchRemoteReader::_fetch_data() { _block = vectorized::Block::create_shared(); st = _block->deserialize(callback->response_->block()); ARROW_RETURN_NOT_OK(to_arrow_status(st)); + LOG(WARNING) << fmt::format( + "ArrowFlightBatchRemoteReader fetch deserialize, finistId={}, packet_seq={}, " + "rows={}, bytes={}", + print_id(_statement->query_id), _packet_seq.load(), _block->rows(), + _block->bytes()); break; } @@ -292,9 +317,19 @@ arrow::Status ArrowFlightBatchRemoteReader::ReadNextImpl(std::shared_ptrquery_id), _packet_seq.load(), _statement->result_addr.hostname, + _statement->result_addr.port, _block != nullptr); ARROW_RETURN_NOT_OK(_fetch_data()); if (_block == nullptr) { // eof, normal path end, last _fetch_data return block is nullptr + LOG(WARNING) << fmt::format( + "ArrowFlightBatchRemoteReader ReadNext eof, finistId={}, packet_seq={}, " + "result_addr={}:{}", + print_id(_statement->query_id), _packet_seq.load(), + _statement->result_addr.hostname, _statement->result_addr.port); return arrow::Status::OK(); } { @@ -308,6 +343,12 @@ arrow::Status ArrowFlightBatchRemoteReader::ReadNextImpl(std::shared_ptrquery_id), _packet_seq.load(), (*out)->num_rows(), + (*out)->num_columns()); VLOG_NOTICE << "ArrowFlightBatchRemoteReader read next: " << (*out)->num_rows() << ", " << (*out)->num_columns() << ", packet_seq: " << _packet_seq; } diff --git a/be/src/service/internal_service.cpp b/be/src/service/internal_service.cpp index daf45179f79914..7e9a7027ca0d12 100644 --- a/be/src/service/internal_service.cpp +++ b/be/src/service/internal_service.cpp @@ -662,14 +662,26 @@ void PInternalService::fetch_arrow_data(google::protobuf::RpcController* control brpc::ClosureGuard closure_guard(done); auto ctx = vectorized::GetArrowResultBatchCtx::create_shared(result); TUniqueId unique_id = UniqueId(request->finst_id()).to_thrift(); // query_id or instance_id + LOG(WARNING) << fmt::format("fetch_arrow_data begin, finistId={}", print_id(unique_id)); std::shared_ptr arrow_buffer; auto st = ExecEnv::GetInstance()->result_mgr()->find_buffer(unique_id, arrow_buffer); if (!st.ok()) { LOG(WARNING) << "Result buffer not found! Query ID: " << print_id(unique_id); return; } + LOG(WARNING) << fmt::format("fetch_arrow_data found buffer, finistId={}, buffer_ptr={}", + print_id(unique_id), fmt::ptr(arrow_buffer.get())); if (st = arrow_buffer->get_batch(ctx); !st.ok()) { LOG(WARNING) << "fetch_arrow_data failed: " << st.to_string(); + } else { + LOG(WARNING) << fmt::format( + "fetch_arrow_data end, finistId={}, response_packet_seq={}, has_eos={}, " + "eos={}, " + "has_empty_batch={}, empty_batch={}, has_block={}, block_bytes={}", + print_id(unique_id), result->has_packet_seq() ? result->packet_seq() : -1, + result->has_eos(), result->has_eos() && result->eos(), + result->has_empty_batch(), result->has_empty_batch() && result->empty_batch(), + result->has_block(), result->has_block() ? result->block().ByteSizeLong() : 0); } }); if (!ret) { diff --git a/be/src/vec/sink/varrow_flight_result_writer.cpp b/be/src/vec/sink/varrow_flight_result_writer.cpp index 9105ba9e057dec..a630b9399771b3 100644 --- a/be/src/vec/sink/varrow_flight_result_writer.cpp +++ b/be/src/vec/sink/varrow_flight_result_writer.cpp @@ -32,14 +32,18 @@ namespace doris::vectorized { void GetArrowResultBatchCtx::on_failure(const Status& status) { DCHECK(!status.ok()) << "status is ok, errmsg=" << status; + LOG(WARNING) << fmt::format("ArrowFlightResult on_failure, status={}", status.to_string()); status.to_protobuf(_result->mutable_status()); } -void GetArrowResultBatchCtx::on_close(int64_t packet_seq, int64_t /* returned_rows */) { +void GetArrowResultBatchCtx::on_close(int64_t packet_seq, int64_t returned_rows) { Status status; status.to_protobuf(_result->mutable_status()); _result->set_packet_seq(packet_seq); _result->set_eos(true); + LOG(WARNING) << fmt::format( + "ArrowFlightResult on_close, packet_seq={}, eos=true, returned_rows={}", packet_seq, + returned_rows); } Status GetArrowResultBatchCtx::on_data(const std::shared_ptr& block, @@ -58,10 +62,17 @@ Status GetArrowResultBatchCtx::on_data(const std::shared_ptr& if (packet_seq == 0) { _result->set_timezone(arrow_buffer->_timezone); } + LOG(WARNING) << fmt::format( + "ArrowFlightResult on_data, finistId={}, packet_seq={}, eos=false, rows={}, " + "block_bytes={}, uncompressed_bytes={}, compressed_bytes={}", + print_id(arrow_buffer->_fragment_id), packet_seq, block->rows(), block->bytes(), + uncompressed_bytes, compressed_bytes); } else { _result->set_empty_batch(true); _result->set_packet_seq(packet_seq); _result->set_eos(false); + LOG(WARNING) << fmt::format( + "ArrowFlightResult on_data empty_batch, packet_seq={}, eos=false", packet_seq); } Status st = Status::OK(); /// The size limit of proto buffer message is 2G