pub struct ShuffleProduceRequest {
pub shuffle_id: u64,
pub side: u8,
pub num_parts: u32,
pub producer_count: u32,
pub keys: Vec<String>,
pub part_node_map: Vec<PartNodeEntry>,
pub plan_bytes: Vec<u8>,
pub tenant_id: u64,
pub database_id: u64,
pub deadline_remaining_ms: u64,
pub trace_id: [u8; 16],
pub descriptor_versions: Vec<DescriptorVersionEntry>,
}Expand description
Cross-node shuffle PRODUCER trigger (E4a).
A coordinator sends this to a producer node to make it execute a LOCAL scan
fragment (plan_bytes), hash-partition each output row on keys, and fan the
rows out to the per-part owners (part_node_map) as ShufflePush streams —
one ShufflePushEnd to EVERY part for this side so each receiver’s per-part
barrier reaches producer_count. The producer replies with exactly one
ShuffleProduceResponse; it never streams the scanned rows back.
side is 0 for the build side and 1 for the probe side of a hash join.
The plan_bytes / tenant_id / database_id / deadline_remaining_ms /
trace_id / descriptor_versions fields mirror
ExecuteRequest so the producer can reuse
the existing local streaming-execution prologue verbatim.
Cross-version safety: new optional fields should be added as Option<T>.
Fields§
§shuffle_id: u64§side: u80 = build side, 1 = probe side — the local side this producer scans.
num_parts: u32§producer_count: u32How many producers will End each (part, side); forwarded into every
emitted ShufflePushRequest so receivers size their build barrier.
keys: Vec<String>Local-side join field names; each output row is hashed on these.
part_node_map: Vec<PartNodeEntry>part -> owning node id (coordinator-computed, sorted by part).
plan_bytes: Vec<u8>Encoded PhysicalPlan of the local scan fragment to execute.
tenant_id: u64§database_id: u64§deadline_remaining_ms: u64§trace_id: [u8; 16]§descriptor_versions: Vec<DescriptorVersionEntry>Trait Implementations§
Source§impl Archive for ShuffleProduceRequest
impl Archive for ShuffleProduceRequest
Source§const COPY_OPTIMIZATION: CopyOptimization<Self>
const COPY_OPTIMIZATION: CopyOptimization<Self>
serialize. Read moreSource§type Archived = ArchivedShuffleProduceRequest
type Archived = ArchivedShuffleProduceRequest
Source§type Resolver = ShuffleProduceRequestResolver
type Resolver = ShuffleProduceRequestResolver
Source§impl Clone for ShuffleProduceRequest
impl Clone for ShuffleProduceRequest
Source§fn clone(&self) -> ShuffleProduceRequest
fn clone(&self) -> ShuffleProduceRequest
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreSource§impl Debug for ShuffleProduceRequest
impl Debug for ShuffleProduceRequest
Source§impl<__D: Fallible + ?Sized> Deserialize<ShuffleProduceRequest, __D> for Archived<ShuffleProduceRequest>where
u64: Archive,
<u64 as Archive>::Archived: Deserialize<u64, __D>,
u8: Archive,
<u8 as Archive>::Archived: Deserialize<u8, __D>,
u32: Archive,
<u32 as Archive>::Archived: Deserialize<u32, __D>,
Vec<String>: Archive,
<Vec<String> as Archive>::Archived: Deserialize<Vec<String>, __D>,
Vec<PartNodeEntry>: Archive,
<Vec<PartNodeEntry> as Archive>::Archived: Deserialize<Vec<PartNodeEntry>, __D>,
Vec<u8>: Archive,
<Vec<u8> as Archive>::Archived: Deserialize<Vec<u8>, __D>,
[u8; 16]: Archive,
<[u8; 16] as Archive>::Archived: Deserialize<[u8; 16], __D>,
Vec<DescriptorVersionEntry>: Archive,
<Vec<DescriptorVersionEntry> as Archive>::Archived: Deserialize<Vec<DescriptorVersionEntry>, __D>,
impl<__D: Fallible + ?Sized> Deserialize<ShuffleProduceRequest, __D> for Archived<ShuffleProduceRequest>where
u64: Archive,
<u64 as Archive>::Archived: Deserialize<u64, __D>,
u8: Archive,
<u8 as Archive>::Archived: Deserialize<u8, __D>,
u32: Archive,
<u32 as Archive>::Archived: Deserialize<u32, __D>,
Vec<String>: Archive,
<Vec<String> as Archive>::Archived: Deserialize<Vec<String>, __D>,
Vec<PartNodeEntry>: Archive,
<Vec<PartNodeEntry> as Archive>::Archived: Deserialize<Vec<PartNodeEntry>, __D>,
Vec<u8>: Archive,
<Vec<u8> as Archive>::Archived: Deserialize<Vec<u8>, __D>,
[u8; 16]: Archive,
<[u8; 16] as Archive>::Archived: Deserialize<[u8; 16], __D>,
Vec<DescriptorVersionEntry>: Archive,
<Vec<DescriptorVersionEntry> as Archive>::Archived: Deserialize<Vec<DescriptorVersionEntry>, __D>,
Source§fn deserialize(
&self,
deserializer: &mut __D,
) -> Result<ShuffleProduceRequest, <__D as Fallible>::Error>
fn deserialize( &self, deserializer: &mut __D, ) -> Result<ShuffleProduceRequest, <__D as Fallible>::Error>
Auto Trait Implementations§
impl Freeze for ShuffleProduceRequest
impl RefUnwindSafe for ShuffleProduceRequest
impl Send for ShuffleProduceRequest
impl Sync for ShuffleProduceRequest
impl Unpin for ShuffleProduceRequest
impl UnsafeUnpin for ShuffleProduceRequest
impl UnwindSafe for ShuffleProduceRequest
Blanket Implementations§
Source§impl<T> ArchivePointee for T
impl<T> ArchivePointee for T
Source§type ArchivedMetadata = ()
type ArchivedMetadata = ()
Source§fn pointer_metadata(
_: &<T as ArchivePointee>::ArchivedMetadata,
) -> <T as Pointee>::Metadata
fn pointer_metadata( _: &<T as ArchivePointee>::ArchivedMetadata, ) -> <T as Pointee>::Metadata
Source§impl<T> ArchiveUnsized for Twhere
T: Archive,
impl<T> ArchiveUnsized for Twhere
T: Archive,
Source§type Archived = <T as Archive>::Archived
type Archived = <T as Archive>::Archived
Archive, it may be
unsized. Read moreSource§fn archived_metadata(
&self,
) -> <<T as ArchiveUnsized>::Archived as ArchivePointee>::ArchivedMetadata
fn archived_metadata( &self, ) -> <<T as ArchiveUnsized>::Archived as ArchivePointee>::ArchivedMetadata
Source§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedExplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
Source§impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
impl<'a, T, E> AsTaggedImplicit<'a, E> for Twhere
T: 'a,
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> LayoutRaw for T
impl<T> LayoutRaw for T
Source§fn layout_raw(_: <T as Pointee>::Metadata) -> Result<Layout, LayoutError>
fn layout_raw(_: <T as Pointee>::Metadata) -> Result<Layout, LayoutError>
Source§impl<T, N1, N2> Niching<NichedOption<T, N1>> for N2
impl<T, N1, N2> Niching<NichedOption<T, N1>> for N2
Source§unsafe fn is_niched(niched: *const NichedOption<T, N1>) -> bool
unsafe fn is_niched(niched: *const NichedOption<T, N1>) -> bool
Source§fn resolve_niched(out: Place<NichedOption<T, N1>>)
fn resolve_niched(out: Place<NichedOption<T, N1>>)
out indicating that a T is niched.impl<T> Read<Exclusive, BecauseExclusive> for Twhere
T: ?Sized,
Source§impl<T, S> SerializeUnsized<S> for T
impl<T, S> SerializeUnsized<S> for T
Source§impl<SS, SP> SupersetOf<SS> for SPwhere
SS: SubsetOf<SP>,
impl<SS, SP> SupersetOf<SS> for SPwhere
SS: SubsetOf<SP>,
Source§fn to_subset(&self) -> Option<SS>
fn to_subset(&self) -> Option<SS>
self from the equivalent element of its
superset. Read moreSource§fn is_in_subset(&self) -> bool
fn is_in_subset(&self) -> bool
self is actually part of its subset T (and can be converted to it).Source§fn to_subset_unchecked(&self) -> SS
fn to_subset_unchecked(&self) -> SS
self.to_subset but without any property checks. Always succeeds.Source§fn from_subset(element: &SS) -> SP
fn from_subset(element: &SS) -> SP
self to the equivalent element of its superset.