Skip to content

Commit 87b234a

Browse files
committed
refactor(flight): own outbound streams directly
1 parent e74473c commit 87b234a

6 files changed

Lines changed: 127 additions & 225 deletions

File tree

src/query/service/src/servers/flight/v1/exchange/exchange_manager.rs

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1874,11 +1874,12 @@ mod tests {
18741874
use databend_common_base::runtime::Runtime;
18751875
use databend_common_exception::ErrorCode;
18761876
use futures::StreamExt;
1877+
use tokio::sync::Semaphore;
18771878

18781879
use super::QueryCoordinator;
18791880
use crate::servers::flight::NewFlightStream;
1880-
use crate::servers::flight::v1::exchange::exchange_packet_sender::ExchangePacketSender;
18811881
use crate::servers::flight::v1::packets::DataPacket;
1882+
use crate::servers::flight::v1::transport::OutboundStream;
18821883
use crate::servers::flight::v1::transport::StreamSendOutcome;
18831884
use crate::servers::flight::v1::transport::frame_lane;
18841885
use crate::servers::flight::v1::transport::reliable::DoExchangeConnector;
@@ -1965,20 +1966,20 @@ mod tests {
19651966
let pending = attached.remove("target").unwrap();
19661967
assert!(attached.is_empty());
19671968

1968-
let sender = ExchangePacketSender::from_reliable(vec![pending], 1, &runtime);
1969+
let stream = pending.start(Arc::new(Semaphore::new(64)), Some(256 * 1024), &runtime);
19691970
let payload = FlightData {
19701971
app_metadata: vec![0, 0].into(),
19711972
data_body: vec![1, 2, 3].into(),
19721973
..Default::default()
19731974
};
19741975
assert_eq!(
1975-
sender.send(0, 0, payload).await.unwrap(),
1976+
stream.send(0, payload).await.unwrap(),
19761977
StreamSendOutcome::Accepted
19771978
);
1978-
sender.finish_producer().await.unwrap();
1979+
stream.finish().await.unwrap();
19791980

19801981
assert!(finish_seen.load(Ordering::SeqCst));
1981-
assert!(sender.destination(0).is_closed());
1982+
assert!(stream.is_closed());
19821983
}
19831984

19841985
#[tokio::test]

src/query/service/src/servers/flight/v1/exchange/exchange_packet_sender.rs

Lines changed: 0 additions & 134 deletions
This file was deleted.

src/query/service/src/servers/flight/v1/exchange/exchange_packet_sink.rs

Lines changed: 8 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -28,25 +28,22 @@ use databend_common_pipeline::sinks::AsyncSink;
2828
use databend_common_pipeline::sinks::AsyncSinker;
2929

3030
use super::serde::ExchangeSerializeMeta;
31-
use crate::servers::flight::v1::exchange::exchange_packet_sender::ExchangePacketSender;
31+
use crate::servers::flight::v1::transport::OutboundStreamRef;
3232
use crate::servers::flight::v1::transport::StreamSendOutcome;
3333

3434
pub struct ExchangePacketSink {
35-
sender: Arc<ExchangePacketSender>,
36-
destination: usize,
35+
stream: OutboundStreamRef,
3736
ignore_exchange: bool,
3837
}
3938

4039
impl ExchangePacketSink {
4140
fn create(
4241
input: Arc<InputPort>,
43-
sender: Arc<ExchangePacketSender>,
44-
destination: usize,
42+
stream: OutboundStreamRef,
4543
ignore_exchange: bool,
4644
) -> ProcessorPtr {
4745
ProcessorPtr::create(AsyncSinker::create(input, Self {
48-
sender,
49-
destination,
46+
stream,
5047
ignore_exchange,
5148
}))
5249
}
@@ -57,7 +54,7 @@ impl AsyncSink for ExchangePacketSink {
5754
const NAME: &'static str = "ExchangePacketSink";
5855

5956
async fn on_finish(&mut self) -> Result<()> {
60-
self.sender.finish_producer().await
57+
self.stream.finish().await
6158
}
6259

6360
async fn consume(&mut self, mut data_block: DataBlock) -> Result<bool> {
@@ -77,9 +74,7 @@ impl AsyncSink for ExchangePacketSink {
7774
bytes += packet.bytes_size();
7875
let flight_data = FlightData::try_from(packet)?;
7976

80-
if self.sender.send(0, self.destination, flight_data).await?
81-
== StreamSendOutcome::ConsumerClosed
82-
{
77+
if self.stream.send(0, flight_data).await? == StreamSendOutcome::ConsumerClosed {
8378
return Ok(true);
8479
}
8580
}
@@ -89,14 +84,10 @@ impl AsyncSink for ExchangePacketSink {
8984
}
9085
}
9186

92-
pub fn create_packet_writer_item(
93-
sender: Arc<ExchangePacketSender>,
94-
destination: usize,
95-
ignore_exchange: bool,
96-
) -> PipeItem {
87+
pub fn create_packet_writer_item(stream: OutboundStreamRef, ignore_exchange: bool) -> PipeItem {
9788
let input = InputPort::create();
9889
PipeItem::create(
99-
ExchangePacketSink::create(input.clone(), sender, destination, ignore_exchange),
90+
ExchangePacketSink::create(input.clone(), stream, ignore_exchange),
10091
vec![input],
10192
vec![],
10293
)

0 commit comments

Comments
 (0)