@@ -27,15 +27,15 @@ let worker = Worker::default()
2727
2828## Querying a worker's version
2929
30- From the coordinating context, use ` DefaultChannelResolver ` to get a cached
31- channel and ` create_worker_client ` to build a client, then call ` get_worker_info ` :
30+ From the coordinating context, use ` grpc:: DefaultChannelResolver` to get a cached
31+ channel and ` grpc:: create_worker_client` to build a client, then call ` get_worker_info ` :
3232
3333``` rust
34- use datafusion_distributed :: {DefaultChannelResolver , GetWorkerInfoRequest , create_worker_client };
34+ use datafusion_distributed :: {grpc , GetWorkerInfoRequest };
3535
36- let channel_resolver = DefaultChannelResolver :: default ();
36+ let channel_resolver = grpc :: DefaultChannelResolver :: default ();
3737let channel = channel_resolver . get_channel (& worker_url ). await ? ;
38- let mut client = create_worker_client (channel );
38+ let mut client = grpc :: create_worker_client (channel );
3939
4040let response = client . get_worker_info (GetWorkerInfoRequest {}). await ? ;
4141println! (" version: {}" , response . into_inner (). version);
@@ -66,9 +66,7 @@ use std::sync::{Arc, RwLock};
6666use std :: time :: Duration ;
6767use url :: Url ;
6868use datafusion :: common :: {HashMap , DataFusionError };
69- use datafusion_distributed :: {
70- DefaultChannelResolver , GetWorkerInfoRequest , WorkerResolver , create_worker_client,
71- };
69+ use datafusion_distributed :: {grpc, GetWorkerInfoRequest , WorkerResolver };
7270
7371struct VersionAwareWorkerResolver {
7472 compatible_urls : Arc <RwLock <Vec <Url >>>,
@@ -80,7 +78,7 @@ async fn background_version_resolver(
8078 all_worker_urls : Vec <Url >,
8179 local_version : String ,
8280 compatible_urls : Arc <RwLock <Vec <Url >>>,
83- channel_resolver : Arc <DefaultChannelResolver >,
81+ channel_resolver : Arc <grpc :: DefaultChannelResolver >,
8482) {
8583 let mut version_cache : HashMap <Url , String > = HashMap :: new ();
8684
@@ -94,7 +92,7 @@ async fn background_version_resolver(
9492 let cr = Arc :: clone (& channel_resolver );
9593 async move {
9694 let channel = cr . get_channel (url ). await . ok ()? ;
97- let mut client = create_worker_client (channel );
95+ let mut client = grpc :: create_worker_client (channel );
9896 let resp = client . get_worker_info (GetWorkerInfoRequest {}). await . ok ()? ;
9997 Some (resp . into_inner (). version)
10098 }
@@ -122,7 +120,7 @@ impl VersionAwareWorkerResolver {
122120 fn start_version_filtering (
123121 all_worker_urls : Vec <Url >,
124122 expected_version : String ,
125- channel_resolver : Arc <DefaultChannelResolver >,
123+ channel_resolver : Arc <grpc :: DefaultChannelResolver >,
126124 ) -> Self {
127125 let compatible_urls = Arc :: new (RwLock :: new (vec! []));
128126
0 commit comments