| 79 | } |
| 80 | |
| 81 | void ClusterFunctionReadTaskResponse::serialize(WriteBuffer & out, size_t worker_protocol_version) const |
| 82 | { |
| 83 | auto protocol_version |
| 84 | = std::min(static_cast<UInt64>(worker_protocol_version), static_cast<UInt64>(DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION)); |
| 85 | writeVarUInt(protocol_version, out); |
| 86 | writeStringBinary(path, out); |
| 87 | |
| 88 | if (protocol_version >= DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_DATA_LAKE_METADATA) |
| 89 | { |
| 90 | SerializedSetsRegistry registry; |
| 91 | if (data_lake_metadata.schema_transform) |
| 92 | data_lake_metadata.schema_transform->serialize(out, registry); |
| 93 | else |
| 94 | ActionsDAG().serialize(out, registry); |
| 95 | |
| 96 | if (protocol_version >= DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_EXCLUDED_ROWS) |
| 97 | { |
| 98 | if (data_lake_metadata.excluded_rows) |
| 99 | data_lake_metadata.excluded_rows->write(out); |
| 100 | else |
| 101 | DataLakeObjectMetadata::ExcludedRows().write(out); |
| 102 | } |
| 103 | } |
| 104 | |
| 105 | if (protocol_version >= DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_FILE_BUCKETS_INFO) |
| 106 | { |
| 107 | if (file_bucket_info) |
| 108 | { |
| 109 | /// Write format name so we can create appropriate file bucket info during deserialization. |
| 110 | writeStringBinary(file_bucket_info->getFormatName(), out); |
| 111 | file_bucket_info->serialize(out); |
| 112 | } |
| 113 | else |
| 114 | { |
| 115 | /// Write empty string as format name if file_bucket_info is not set. |
| 116 | writeStringBinary("", out); |
| 117 | } |
| 118 | } |
| 119 | |
| 120 | if (protocol_version >= DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION_WITH_ICEBERG_METADATA) |
| 121 | { |
| 122 | if (iceberg_info.has_value()) |
| 123 | { |
| 124 | writeVarUInt(1, out); |
| 125 | iceberg_info->serializeForClusterFunctionProtocol(out, protocol_version); |
| 126 | } |
| 127 | else |
| 128 | { |
| 129 | writeVarUInt(0, out); |
| 130 | } |
| 131 | } |
| 132 | } |
| 133 | |
| 134 | void ClusterFunctionReadTaskResponse::deserialize(ReadBuffer & in) |
| 135 | { |
nothing calls this directly
no test coverage detected