mrhhsg commented on code in PR #63258:
URL: https://github.com/apache/doris/pull/63258#discussion_r4229106903


##########
be/src/util/proto_util.h:
##########
@@ -39,9 +40,24 @@ constexpr size_t MIN_HTTP_BRPC_SIZE = (1ULL << 31);
 // Embed column_values and brpc request serialization string in controller 
attachment.
 template <typename Params, typename Closure>
 Status request_embed_attachment_contain_blockv2(Params* brpc_request,
-                                                std::unique_ptr<Closure>& 
closure) {
-    std::string column_values = 
std::move(*brpc_request->mutable_block()->mutable_column_values());
-    brpc_request->mutable_block()->mutable_column_values()->clear();
+                                                std::unique_ptr<Closure>& 
closure,
+                                                bool restore_column_values = 
false) {
+    auto* block = brpc_request->mutable_block();
+    if (restore_column_values) {
+        // Some callers borrow block storage from a shared owner. Temporarily 
detach the large
+        // column_values field so the serialized request stays small, then 
restore it before
+        // returning so the real owner can still be reused by later sends.
+        auto* column_values = block->release_column_values();
+        DORIS_CHECK(column_values != nullptr);

Review Comment:
   The restore mode is gone. Since 27e530b3639 the broadcast HTTP path no 
longer touches the shared `PBlock`: it builds a metadata-only request copy 
(`make_http_request_without_column_values`) and passes the holder's 
`column_values` as read-only attachment data through 
`transmit_block_httpv2_with_attachment_data`, so concurrent sends of the same 
holder from different `RpcInstance` callbacks only read it. Resolving as 
addressed.



##########
be/src/storage/rowset/rowset_meta.cpp:
##########
@@ -300,6 +301,9 @@ bool RowsetMeta::_deserialize_from_pb(std::string_view 
value) {
                     _rowset_meta_pb.inverted_index_storage_format());
         }
         set_tablet_schema(schema_pb);
+        // The schema has been materialized into TabletSchemaCache by 
set_tablet_schema(). Drop the
+        // protobuf-owned copy from `_rowset_meta_pb` to avoid holding the 
large schema twice; passing
+        // nullptr intentionally deletes the current protobuf submessage.
         _rowset_meta_pb.set_allocated_tablet_schema(nullptr);

Review Comment:
   Not changed in this PR. `_deserialize_from_pb()` clearing `tablet_schema` 
with `set_allocated_tablet_schema(nullptr)` and `json_rowset_meta()` 
serializing `_rowset_meta_pb` directly both already exist on master (merge base 
9a5f86ec82b, `rowset_meta.cpp:303`); this PR only adds a comment above that 
line. `json_rowset_meta()` also has no production caller under `be/src`. The 
JSON export behaviour is pre-existing and out of scope for this exchange fix, 
so I am leaving it as is here.



##########
be/src/exec/operator/exchange_sink_buffer.cpp:
##########
@@ -49,6 +49,49 @@
 
 namespace doris {
 
+namespace exchange_sink_buffer::detail {
+
+void copy_block_metadata_without_column_values(const PBlock& src, PBlock* dst) 
{
+    for (int i = 0; i < src.column_metas_size(); ++i) {
+        dst->add_column_metas()->CopyFrom(src.column_metas(i));
+    }
+    dst->set_be_exec_version(src.be_exec_version());
+    dst->set_compressed(src.compressed());
+    dst->set_compression_type(src.compression_type());
+    dst->set_uncompressed_size(src.uncompressed_size());
+}
+
+std::shared_ptr<PTransmitDataParams> make_http_request_without_column_values(
+        const PTransmitDataParams& src) {
+    DORIS_CHECK(src.has_block());
+    DORIS_CHECK(src.blocks_size() == 0);
+    DORIS_CHECK(!src.has_row_batch());
+
+    auto dst = std::make_shared<PTransmitDataParams>();
+    dst->mutable_finst_id()->CopyFrom(src.finst_id());
+    dst->set_node_id(src.node_id());
+    dst->set_sender_id(src.sender_id());
+    dst->set_be_number(src.be_number());
+    dst->set_eos(src.eos());
+    dst->set_packet_seq(src.packet_seq());
+    if (src.has_query_statistics()) {
+        dst->mutable_query_statistics()->CopyFrom(src.query_statistics());
+    }
+    if (src.has_transfer_by_attachment()) {
+        dst->set_transfer_by_attachment(src.transfer_by_attachment());
+    }
+    if (src.has_query_id()) {
+        dst->mutable_query_id()->CopyFrom(src.query_id());
+    }
+    if (src.has_exec_status()) {
+        dst->mutable_exec_status()->CopyFrom(src.exec_status());
+    }
+    copy_block_metadata_without_column_values(src.block(), 
dst->mutable_block());

Review Comment:
   Fixed in 8d44f49991e. `copy_block_metadata_without_column_values()` now 
marks `column_values` present with an empty value 
(`dst->set_column_values("")`), the same shape 
`request_embed_attachment_contain_blockv2()` leaves behind after moving the 
owned payload out, so the receiver's `ParseFromString()` sets the presence bit 
and the `PBlock::CopyFrom()` in `VDataStreamMgr::transmit_block()` keeps the 
extracted bytes.
   
   Tests: added 
`ExchangeSinkBufferTest.HttpRequestKeepsColumnValuesThroughAttachmentAndCopyFrom`
 (metadata-only request + `request_embed_attachmentv2` with the borrowed 
payload -> `attachment_extract_request_contain_block` -> `PBlock::CopyFrom`) 
and `ProtoUtilTest.AttachmentWithoutColumnValuesPresenceIsDroppedByCopyFrom`, 
which documents the failure shape; the existing tests now assert presence after 
extraction and after `CopyFrom`.
   
   Verified end to end on a local 1 FE + 2 BE cluster with the forced-HTTP 
debug point: without this line the broadcast join fails on the receiver 
(`Allocator sys memory check failed: Cannot alloc:134217728.00 GB`, from 
deserializing an empty payload); with it both queries pass.



##########
regression-test/suites/correctness_p0/test_exchange_http_send_borrowed_block.groovy:
##########
@@ -0,0 +1,74 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+import org.apache.doris.regression.suite.ClusterOptions
+
+// Forces exchange sinks to use the http-attachment path (normally reserved for
+// requests >= 2G) via a debug point, so that both the unicast (shuffle) and
+// broadcast exchange sink code in proto_util.h/exchange_sink_buffer.cpp that
+// move/borrow column_values into the brpc attachment actually run cross-BE,
+// instead of only through unit tests.
+suite('test_exchange_http_send_borrowed_block', 'docker') {
+    def options = new ClusterOptions()
+    options.beNum = 2
+    options.enableDebugPoints()
+
+    docker(options) {
+        def tbl = 'test_exchange_http_send_borrowed_block_tbl'
+
+        sql "DROP TABLE IF EXISTS ${tbl}"
+        sql """
+            CREATE TABLE ${tbl} (
+                k1 INT NULL,
+                k2 INT NULL
+            )
+            DISTRIBUTED BY HASH(k1) BUCKETS 4
+            PROPERTIES ("replication_num" = "1")
+        """
+        sql """
+            INSERT INTO ${tbl} VALUES
+            (1, 10), (2, 20), (3, 30), (4, 40), (5, 50), (6, 60), (7, 70), (8, 
80)
+        """
+
+        def backends = cluster.getAllBackends()
+        assertTrue(backends.size() == 2)
+        for (def be : backends) {
+            
be.enableDebugPoint('proto_util.enable_http_send_block.always_http', null)
+        }
+
+        try {
+            // Shuffle exchange: exercises the unicast attachment path

Review Comment:
   Fixed in 8d44f49991e. The suite now runs `SET 
exchange_multi_blocks_byte_size = -1` before the queries, and after each query 
asserts that the sum over all BEs of brpc's built-in 
`rpc_server_<port>_doris_pbackend_service_transmit_block_by_http_count` (read 
from `/brpc_metrics` via `WarmupMetricsUtils.getBrpcMetric`) increased, so a 
run cannot pass without reaching `transmit_block_by_http`.
   
   While verifying this I also found that the self join on the distribution 
column `k1` was planned as a COLOCATE join with no exchange at all, so both 
queries now join on `k2` (same result rows, `.out` unchanged). On a local 2 BE 
cluster the counter goes 0 -> 6 after the shuffle join and -> 10 after the 
broadcast join, and the broadcast query fails there without the presence fix 
above.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to