risingwave_stream/from_proto/iceberg_with_pk_index/
writer.rs1use std::sync::Arc;
16
17use anyhow::anyhow;
18use risingwave_common::secret::LocalSecretManager;
19use risingwave_connector::sink::iceberg::{
20 ICEBERG_SINK, IcebergConfig, create_and_validate_table_impl,
21};
22use risingwave_connector::sink::{SinkMetaClient, SinkWriterParam};
23use risingwave_expr::bail;
24use risingwave_pb::id::SinkId;
25use risingwave_pb::stream_plan::IcebergWithPkIndexWriterNode;
26use risingwave_storage::StateStore;
27
28use super::super::sink::build_sink_param;
29use crate::common::table::state_table::{StateTableBuilder, StateTableOpConsistencyLevel};
30use crate::error::StreamResult;
31use crate::executor::{Executor, IcebergWriterImpl, StreamExecutorError, WriterExecutor};
32use crate::from_proto::ExecutorBuilder;
33use crate::task::ExecutorParams;
34pub struct IcebergWithPkIndexWriterExecutorBuilder;
35
36impl_stream_node_body!(IcebergWithPkIndexWriter(IcebergWithPkIndexWriterNode) => IcebergWithPkIndexWriterExecutorBuilder);
37
38impl ExecutorBuilder for IcebergWithPkIndexWriterExecutorBuilder {
39 type Node = IcebergWithPkIndexWriterNode;
40
41 async fn new_boxed_executor(
42 params: ExecutorParams,
43 node: &Self::Node,
44 store: impl StateStore,
45 ) -> StreamResult<Executor> {
46 let [input, resolver_input] = params
47 .input
48 .try_into()
49 .map_err(|_| anyhow!("IcebergWithPkIndexWriterExecutor requires exactly two inputs"))?;
50 let sink_desc = node.sink_desc.as_ref().unwrap();
51 let sink_id: SinkId = sink_desc.get_id();
52 let sink_name = sink_desc.get_name().to_owned();
53
54 let properties_with_secret = LocalSecretManager::global().fill_secrets(
55 sink_desc.get_properties().clone(),
56 sink_desc.get_secret_refs().clone(),
57 )?;
58 let config = IcebergConfig::from_btreemap(properties_with_secret.clone())
59 .map_err(|err| StreamExecutorError::from((err, sink_id)))?;
60
61 let pk_indices = sink_desc
62 .downstream_pk
63 .iter()
64 .map(|&idx| idx as usize)
65 .collect::<Vec<_>>();
66 if pk_indices.is_empty() {
67 bail!("missing downstream pk in iceberg sink desc");
68 }
69
70 let (sink_param, _) = build_sink_param(sink_desc, properties_with_secret, ICEBERG_SINK)?;
71
72 let table = create_and_validate_table_impl(&config, &sink_param)
73 .await
74 .map_err(|e| StreamExecutorError::sink_error(e, sink_id))?;
75
76 let pk_index_state_table = StateTableBuilder::new(
77 node.get_pk_index_table()?,
78 store,
79 params.vnode_bitmap.clone().map(Arc::new),
80 )
81 .enable_preload_all_rows_by_config(¶ms.config)
82 .with_op_consistency_level(StateTableOpConsistencyLevel::Inconsistent)
83 .build()
84 .await;
85
86 let meta_client = params
87 .env
88 .meta_client()
89 .ok_or_else(|| anyhow!("meta client is required for Iceberg writer"))?;
90 let meta_client = SinkMetaClient::MetaClient(meta_client);
91
92 let writer_param = SinkWriterParam {
93 executor_id: params.executor_id,
94 vnode_bitmap: params.vnode_bitmap.clone(),
95 meta_client: Some(meta_client),
96 extra_partition_col_idx: sink_desc.extra_partition_col_idx.map(|v| v as usize),
97 actor_id: params.actor_context.id,
98 sink_id,
99 sink_name,
100 connector: ICEBERG_SINK.to_owned(),
101 streaming_config: params.config.as_ref().clone(),
102 time_zone: params.actor_context.time_zone,
103 };
104 let writer = IcebergWriterImpl::build(&config, table, &writer_param)?;
105
106 let exec = WriterExecutor::new(
107 params.actor_context,
108 input,
109 resolver_input,
110 pk_indices,
111 pk_index_state_table,
112 writer,
113 params.config.developer.chunk_size,
114 sink_id,
115 params.local_barrier_manager.clone(),
116 );
117 Ok((params.info, exec).into())
118 }
119}