Skip to content

Commit b25ff06

Browse files
committed
feat(query): persist captured lineage
1 parent 24e634c commit b25ff06

18 files changed

Lines changed: 715 additions & 1 deletion

src/query/service/src/interpreters/common/lineage_writer.rs

Lines changed: 605 additions & 0 deletions
Large diffs are not rendered by default.

src/query/service/src/interpreters/common/mod.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
mod column;
1616
mod finish_hook;
1717
mod grant;
18+
mod lineage_writer;
1819
mod metrics;
1920
mod notification;
2021
mod query_log;
@@ -29,6 +30,9 @@ pub mod table_option_validation;
2930
pub use column::*;
3031
pub use finish_hook::QueryFinishHooks;
3132
pub use grant::validate_grant_object_exists;
33+
pub use lineage_writer::attach_query_lineage_on_finished;
34+
pub use lineage_writer::build_query_lineage_updates;
35+
pub use lineage_writer::build_query_lineage_updates_with_lineage;
3236
pub use log::*;
3337
pub use notification::get_notification_client_config;
3438
pub use query_log::InterpreterQueryLog;

src/query/service/src/interpreters/interpreter_copy_into_table.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,8 @@ use log::info;
6767
use crate::interpreters::HookOperator;
6868
use crate::interpreters::Interpreter;
6969
use crate::interpreters::SelectInterpreter;
70+
use crate::interpreters::common::attach_query_lineage_on_finished;
71+
use crate::interpreters::common::build_query_lineage_updates;
7072
use crate::interpreters::common::check_deduplicate_label;
7173
use crate::interpreters::common::dml_build_update_stream_req;
7274
use crate::physical_plans::CopyIntoTable;
@@ -874,6 +876,7 @@ impl CopyIntoTableInterpreter {
874876
path_prefix,
875877
)?;
876878

879+
let lineage_updates = build_query_lineage_updates(&ctx, None).await?;
877880
to_table.commit_insertion(
878881
ctx.clone(),
879882
main_pipeline,
@@ -884,6 +887,7 @@ impl CopyIntoTableInterpreter {
884887
deduplicated_label,
885888
table_meta_timestamps,
886889
)?;
890+
attach_query_lineage_on_finished(main_pipeline, lineage_updates);
887891
}
888892

889893
// Purge files.

src/query/service/src/interpreters/interpreter_factory.rs

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ use std::sync::Arc;
1616

1717
use databend_common_ast::ast::ExplainKind;
1818
use databend_common_catalog::lock::LockTableOption;
19+
use databend_common_config::GlobalConfig;
1920
use databend_common_exception::ErrorCode;
2021
use databend_common_exception::Result;
2122
use databend_common_sql::binder::ExplainConfig;
@@ -153,6 +154,7 @@ impl InterpreterFactory {
153154
error!("Access.denied(v2): {:?}", e);
154155
}
155156
})?;
157+
initialize_query_lineage(&ctx, plan);
156158
let mut access_logger = AccessLogger::create(ctx.clone());
157159
access_logger.log(plan);
158160
access_logger.output();
@@ -937,3 +939,24 @@ impl InterpreterFactory {
937939
}
938940
}
939941
}
942+
943+
fn initialize_query_lineage(ctx: &QueryContext, plan: &Plan) {
944+
ctx.init_query_lineage(|| {
945+
if !GlobalConfig::instance()
946+
.query
947+
.common
948+
.lineage
949+
.capture_enabled
950+
{
951+
return None;
952+
}
953+
954+
match plan.query_lineage() {
955+
Ok(lineage) => lineage,
956+
Err(error) => {
957+
log::warn!("failed to extract query lineage: {error:?}");
958+
None
959+
}
960+
}
961+
});
962+
}

src/query/service/src/interpreters/interpreter_insert.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,8 @@ use crate::clusters::ClusterHelper;
6161
use crate::interpreters::HookOperator;
6262
use crate::interpreters::Interpreter;
6363
use crate::interpreters::InterpreterPtr;
64+
use crate::interpreters::common::attach_query_lineage_on_finished;
65+
use crate::interpreters::common::build_query_lineage_updates;
6466
use crate::interpreters::common::check_deduplicate_label;
6567
use crate::interpreters::common::dml_build_update_stream_req;
6668
use crate::physical_plans::ConstantTableScan;
@@ -468,6 +470,8 @@ impl Interpreter for InsertInterpreter {
468470
build_query_pipeline_without_render_result_set(&self.ctx, &insert_select_plan)
469471
.await?;
470472

473+
let lineage_updates =
474+
build_query_lineage_updates(&self.ctx, Some(table.get_table_info())).await?;
471475
table.commit_insertion(
472476
self.ctx.clone(),
473477
&mut build_res.main_pipeline,
@@ -478,6 +482,7 @@ impl Interpreter for InsertInterpreter {
478482
unsafe { self.ctx.get_settings().get_deduplicate_label()? },
479483
table_meta_timestamps,
480484
)?;
485+
attach_query_lineage_on_finished(&mut build_res.main_pipeline, lineage_updates);
481486

482487
// Execute the hook operator.
483488
if self.plan.branch.is_none() {

src/query/service/src/interpreters/interpreter_insert_multi_table.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,8 @@ use databend_common_storages_fuse::FuseTable;
4444
use super::HookOperator;
4545
use crate::interpreters::Interpreter;
4646
use crate::interpreters::InterpreterPtr;
47+
use crate::interpreters::common::attach_query_lineage_on_finished;
48+
use crate::interpreters::common::build_query_lineage_updates;
4749
use crate::interpreters::common::dml_build_update_stream_req;
4850
use crate::physical_plans::CastSchema;
4951
use crate::physical_plans::ChunkAppendData;
@@ -107,6 +109,9 @@ impl Interpreter for InsertMultiTableInterpreter {
107109
let physical_plan = self.build_physical_plan(false).await?;
108110
let mut build_res =
109111
build_query_pipeline_without_render_result_set(&self.ctx, &physical_plan).await?;
112+
let lineage_updates = build_query_lineage_updates(&self.ctx, None).await?;
113+
attach_query_lineage_on_finished(&mut build_res.main_pipeline, lineage_updates);
114+
110115
// Execute hook.
111116
if self
112117
.ctx

src/query/service/src/interpreters/interpreter_mutation.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -178,6 +178,7 @@ impl MutationInterpreter {
178178
let mut builder =
179179
PhysicalPlanBuilder::new(mutation.metadata.clone(), self.ctx.clone(), dry_run);
180180
builder.set_mutation_build_info(mutation_build_info);
181+
builder.set_query_lineage(self.ctx.get_query_lineage());
181182
builder
182183
.build(&self.s_expr, *mutation.required_columns.clone())
183184
.await

src/query/service/src/interpreters/interpreter_replace.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,7 @@ use crate::interpreters::HookOperator;
4747
use crate::interpreters::Interpreter;
4848
use crate::interpreters::InterpreterPtr;
4949
use crate::interpreters::SelectInterpreter;
50+
use crate::interpreters::common::build_query_lineage_updates;
5051
use crate::interpreters::common::check_deduplicate_label;
5152
use crate::interpreters::common::dml_build_update_stream_req;
5253
#[cfg(feature = "storage-stage")]
@@ -358,6 +359,7 @@ impl ReplaceInterpreter {
358359
});
359360
}
360361

362+
let lineage_updates = build_query_lineage_updates(&self.ctx, None).await?;
361363
root = PhysicalPlan::new(CommitSink {
362364
input: root,
363365
snapshot: base_snapshot,
@@ -370,6 +372,7 @@ impl ReplaceInterpreter {
370372
deduplicated_label: unsafe { self.ctx.get_settings().get_deduplicate_label()? },
371373
table_meta_timestamps,
372374
recluster_info: None,
375+
lineage_updates,
373376
meta: PhysicalPlanMeta::new("CommitSink"),
374377
});
375378

src/query/service/src/interpreters/interpreter_table_create.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@ use log::info;
6969

7070
use crate::interpreters::InsertInterpreter;
7171
use crate::interpreters::Interpreter;
72+
use crate::interpreters::common::build_query_lineage_updates;
7273
use crate::interpreters::common::table_option_validation::is_valid_analyze_count_min_sketch_error_rate;
7374
use crate::interpreters::common::table_option_validation::is_valid_analyze_frequency_columns;
7475
use crate::interpreters::common::table_option_validation::is_valid_analyze_histogram_algorithm;
@@ -240,6 +241,7 @@ impl CreateTableInterpreter {
240241
TableIdent::new(table_id, table_id_seq),
241242
table_meta,
242243
);
244+
let lineage_updates = build_query_lineage_updates(&self.ctx, Some(&table_info)).await?;
243245

244246
let insert_plan = Insert {
245247
catalog: self.plan.catalog.clone(),
@@ -288,6 +290,7 @@ impl CreateTableInterpreter {
288290
"create_table_as_select {} success, commit table meta data by table id {}",
289291
qualified_table_name, table_id
290292
);
293+
let lineage_updates = lineage_updates.clone();
291294
let fut = async move {
292295
let req = CommitTableMetaReq {
293296
name_ident: TableNameIdent {
@@ -299,6 +302,7 @@ impl CreateTableInterpreter {
299302
table_id,
300303
prev_table_id,
301304
orphan_table_name,
305+
lineage_updates,
302306
};
303307
catalog.commit_table_meta(req).await
304308
};
@@ -536,6 +540,7 @@ impl CreateTableInterpreter {
536540
table_name: self.plan.table.to_string(),
537541
},
538542
table_meta,
543+
lineage_updates: vec![],
539544
as_dropped: false,
540545
table_properties: self.plan.table_properties.clone(),
541546
table_partition: self.plan.table_partition.as_ref().map(|table_partition| {

src/query/service/src/interpreters/interpreter_table_recluster.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -440,6 +440,7 @@ impl ReclusterTableInterpreter {
440440
deduplicated_label: None,
441441
table_meta_timestamps,
442442
recluster_info,
443+
lineage_updates: vec![],
443444
meta: PhysicalPlanMeta::new("CommitSink"),
444445
})
445446
}

0 commit comments

Comments
 (0)