@@ -7,7 +7,7 @@ use tracing::{error, info};
77
88use sketch_core:: config:: { self , ImplMode } ;
99
10- use query_engine_rust:: data_model:: enums:: { InputFormat , LockStrategy , StreamingEngine } ;
10+ use query_engine_rust:: data_model:: enums:: { InputFormat , LockStrategy , QueryProtocol , StreamingEngine } ;
1111use query_engine_rust:: drivers:: AdapterConfig ;
1212use query_engine_rust:: utils:: file_io:: { read_inference_config, read_streaming_config} ;
1313use query_engine_rust:: {
@@ -233,11 +233,24 @@ async fn main() -> Result<()> {
233233 ) ) ;
234234
235235 // Setup Kafka consumer (equivalent to Python's kafka_thread)
236+ // When deleting the DB, use a unique group ID so no prior committed offset exists
237+ // and the consumer replays from the beginning of the topic.
238+ let kafka_group_id = if args. delete_existing_db {
239+ format ! (
240+ "query-engine-rust-{}" ,
241+ std:: time:: SystemTime :: now( )
242+ . duration_since( std:: time:: UNIX_EPOCH )
243+ . unwrap_or_default( )
244+ . as_secs( )
245+ )
246+ } else {
247+ "query-engine-rust" . to_string ( )
248+ } ;
236249 let kafka_config = KafkaConsumerConfig {
237250 broker : args. kafka_broker . clone ( ) ,
238251 topic : args. kafka_topic . clone ( ) ,
239- group_id : "query-engine-rust" . to_string ( ) ,
240- auto_offset_reset : "beginning " . to_string ( ) ,
252+ group_id : kafka_group_id ,
253+ auto_offset_reset : "earliest " . to_string ( ) ,
241254 input_format : args. input_format ,
242255 decompress_json : args. decompress_json ,
243256 batch_size : 1000 ,
@@ -328,11 +341,22 @@ async fn main() -> Result<()> {
328341 // true, // Always forward (fallback for every query)
329342 //);
330343
331- // Original Prometheus config (commented out temporarily):
332- let adapter_config = AdapterConfig :: prometheus_promql (
333- args. prometheus_server . clone ( ) ,
334- args. forward_unsupported_queries ,
335- ) ;
344+ let adapter_config = match args. query_language {
345+ QueryLanguage :: sql => AdapterConfig :: new (
346+ QueryProtocol :: ClickHouseHttp ,
347+ QueryLanguage :: sql,
348+ None ,
349+ ) ,
350+ QueryLanguage :: elastic_sql => AdapterConfig :: new (
351+ QueryProtocol :: ElasticHttp ,
352+ QueryLanguage :: elastic_sql,
353+ None ,
354+ ) ,
355+ _ => AdapterConfig :: prometheus_promql (
356+ args. prometheus_server . clone ( ) ,
357+ args. forward_unsupported_queries ,
358+ ) ,
359+ } ;
336360
337361 let http_config = HttpServerConfig {
338362 port : args. http_port ,
0 commit comments