risingwave_stream/from_proto/iceberg_with_pk_index/
position_delete_merger.rs1use anyhow::anyhow;
16use risingwave_common::secret::LocalSecretManager;
17use risingwave_connector::sink::iceberg::IcebergConfig;
18use risingwave_pb::id::SinkId;
19use risingwave_pb::stream_plan::IcebergWithPkIndexPositionDeleteMergerNode;
20use risingwave_storage::StateStore;
21
22use crate::error::StreamResult;
23use crate::executor::{
24 Executor, PositionDeleteHandlerImpl, PositionDeleteMergerExecutor, StreamExecutorError,
25};
26use crate::from_proto::ExecutorBuilder;
27use crate::task::ExecutorParams;
28
29pub struct IcebergWithPkIndexPositionDeleteMergerExecutorBuilder;
30
31impl_stream_node_body!(IcebergWithPkIndexPositionDeleteMerger(IcebergWithPkIndexPositionDeleteMergerNode) => IcebergWithPkIndexPositionDeleteMergerExecutorBuilder);
32
33impl ExecutorBuilder for IcebergWithPkIndexPositionDeleteMergerExecutorBuilder {
34 type Node = IcebergWithPkIndexPositionDeleteMergerNode;
35
36 async fn new_boxed_executor(
37 params: ExecutorParams,
38 node: &Self::Node,
39 _store: impl StateStore,
40 ) -> StreamResult<Executor> {
41 let [input]: [_; 1] = params.input.try_into().unwrap();
42
43 let sink_desc = node.sink_desc.as_ref().unwrap();
44 let sink_id: SinkId = sink_desc.get_id();
45
46 let properties_with_secret = LocalSecretManager::global().fill_secrets(
47 sink_desc.get_properties().clone(),
48 sink_desc.get_secret_refs().clone(),
49 )?;
50 let config = IcebergConfig::from_btreemap(properties_with_secret)
51 .map_err(|err| StreamExecutorError::sink_error(err, sink_id))?;
52 let meta_client = params.env.meta_client().ok_or_else(|| {
53 anyhow!("meta client is required for iceberg pk-index position-delete merger")
54 })?;
55
56 let handler = PositionDeleteHandlerImpl::new(
57 config,
58 params.vnode_bitmap.clone(),
59 sink_id,
60 meta_client,
61 );
62
63 let exec = PositionDeleteMergerExecutor::new(
64 params.actor_context.id,
65 sink_id,
66 params.local_barrier_manager.clone(),
67 input,
68 handler,
69 );
70 Ok((params.info, exec).into())
71 }
72}