MCPcopy Create free account
hub / github.com/ClickHouse/ClickHouse / processPacket

Method processPacket

src/QueryPipeline/RemoteQueryExecutor.cpp:637–738  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

635}
636
637RemoteQueryExecutor::ReadResult RemoteQueryExecutor::processPacket(Packet packet)
638{
639 switch (packet.type)
640 {
641 case Protocol::Server::MergeTreeReadTaskRequest:
642 chassert(packet.request.has_value());
643 processMergeTreeReadTaskRequest(packet.request.value());
644 return ReadResult(ReadResult::Type::ParallelReplicasToken);
645
646 case Protocol::Server::MergeTreeAllRangesAnnouncement:
647 chassert(packet.announcement.has_value());
648 processMergeTreeInitialReadAnnouncement(packet.announcement.value());
649 return ReadResult(ReadResult::Type::ParallelReplicasToken);
650
651 case Protocol::Server::ReadTaskRequest:
652 processReadTaskRequest();
653 break;
654 case Protocol::Server::PartUUIDs:
655 LOG_WARNING(
656 log,
657 "The remote server has sent no longer supported packet (Server::PartUUIDs). allow_experimental_query_deduplication feature "
658 "has been deprecated. Consider upgrading the remote server ({})",
659 connections->dumpAddresses());
660 break;
661 case Protocol::Server::Data:
662 /// Note: `packet.block.rows() > 0` means it's a header block.
663 /// We can actually return it, and the first call to RemoteQueryExecutor::read
664 /// will return earlier. We should consider doing it.
665 if (!packet.block.empty() && (packet.block.rows() > 0))
666 return ReadResult(adaptBlockStructure(packet.block, *header));
667 break; /// If the block is empty - we will receive other packets before EndOfStream.
668
669 case Protocol::Server::Exception:
670 got_exception_from_replica = true;
671 packet.exception->rethrow();
672 break;
673
674 case Protocol::Server::EndOfStream:
675 if (!connections->hasActiveConnections())
676 {
677 finished = true;
678 /// TODO: Replace with Type::Finished
679 return ReadResult(Block{});
680 }
681 break;
682
683 case Protocol::Server::Progress:
684 /** We use the progress from a remote server.
685 * We also include in ProcessList,
686 * and we use it to check
687 * constraints (for example, the minimum speed of query execution)
688 * and quotas (for example, the number of lines to read).
689 */
690 if (progress_callback)
691 progress_callback(packet.progress);
692 break;
693
694 case Protocol::Server::ProfileInfo:

Callers

nothing calls this directly

Calls 11

adaptBlockStructureFunction · 0.85
pushBlockMethod · 0.80
ReadResultClass · 0.70
ExceptionClass · 0.50
valueMethod · 0.45
dumpAddressesMethod · 0.45
emptyMethod · 0.45
rowsMethod · 0.45
rethrowMethod · 0.45
hasActiveConnectionsMethod · 0.45
emplaceMethod · 0.45

Tested by

no test coverage detected