1#![expect(
18 clippy::disallowed_types,
19 reason = "construct iceberg::Error to implement the trait"
20)]
21
22use std::collections::HashMap;
23use std::fmt::Debug;
24use std::sync::Arc;
25
26use anyhow::Context;
27use async_trait::async_trait;
28use iceberg::io::FileIOBuilder;
29use iceberg::spec::{Schema, SortOrder, TableMetadata, UnboundPartitionSpec};
30use iceberg::table::Table;
31use iceberg::{
32 Catalog, Namespace, NamespaceIdent, Runtime, TableCommit, TableCreation, TableIdent,
33 TableRequirement, TableUpdate,
34};
35use iceberg_storage_opendal::OpenDalResolvingStorageFactory;
36use itertools::Itertools;
37use jni::objects::{GlobalRef, JObject};
38use risingwave_common::global_jvm::Jvm;
39use risingwave_jni_core::call_method;
40use risingwave_jni_core::jvm_runtime::{execute_with_jni_env, jobj_to_str};
41use serde::{Deserialize, Serialize};
42use thiserror_ext::AsReport;
43
44use crate::error::ConnectorResult;
45
46#[derive(Debug, Deserialize)]
47#[serde(rename_all = "kebab-case")]
48struct LoadTableResponse {
49 pub metadata_location: Option<String>,
50 pub metadata: TableMetadata,
51 pub _config: Option<HashMap<String, String>>,
52}
53
54#[derive(Debug, Serialize)]
55#[serde(rename_all = "kebab-case")]
56struct CreateTableRequest {
57 pub name: String,
59 pub location: Option<String>,
61 pub schema: Schema,
63 pub partition_spec: Option<UnboundPartitionSpec>,
65 pub write_order: Option<SortOrder>,
67 pub properties: HashMap<String, String>,
69}
70
71#[derive(Debug, Serialize, Deserialize)]
72struct CommitTableRequest {
73 identifier: TableIdent,
74 requirements: Vec<TableRequirement>,
75 updates: Vec<TableUpdate>,
76}
77
78#[derive(Debug, Serialize, Deserialize)]
79#[serde(rename_all = "kebab-case")]
80struct CommitTableResponse {
81 metadata_location: String,
82 metadata: TableMetadata,
83}
84
85#[derive(Debug, Serialize, Deserialize)]
86#[serde(rename_all = "kebab-case")]
87struct ListNamespacesResponse {
88 namespaces: Vec<NamespaceIdent>,
89 next_page_token: Option<String>,
90}
91
92#[derive(Debug, Deserialize)]
93struct GetNamespaceResponse {
94 namespace: NamespaceIdent,
95 #[serde(default)]
96 properties: HashMap<String, String>,
97}
98
99#[derive(Debug, Serialize, Deserialize)]
100#[serde(rename_all = "kebab-case")]
101struct ListTablesResponse {
102 identifiers: Vec<TableIdent>,
103 next_page_token: Option<String>,
104}
105
106impl From<&TableCreation> for CreateTableRequest {
107 fn from(value: &TableCreation) -> Self {
108 Self {
109 name: value.name.clone(),
110 location: value.location.clone(),
111 schema: value.schema.clone(),
112 partition_spec: value.partition_spec.clone(),
113 write_order: value.sort_order.clone(),
114 properties: value.properties.clone(),
115 }
116 }
117}
118
119fn namespace_to_string(namespace: &NamespaceIdent) -> String {
120 namespace.iter().join(".")
121}
122
123#[derive(Debug)]
124struct JniCatalogInner {
125 java_catalog: GlobalRef,
126 jvm: Jvm,
127}
128
129#[derive(Debug)]
130pub struct JniCatalog {
131 inner: Arc<JniCatalogInner>,
134 file_io_props: Arc<HashMap<String, String>>,
135}
136
137async fn execute_blocking_jni<T>(
141 task: impl FnOnce() -> anyhow::Result<T> + Send + 'static,
142) -> anyhow::Result<T>
143where
144 T: Send + 'static,
145{
146 tokio::task::spawn_blocking(task)
147 .await
148 .context("Failed to join blocking Iceberg JNI catalog task")?
149}
150
151#[async_trait]
152impl Catalog for JniCatalog {
153 async fn list_namespaces(
155 &self,
156 _parent: Option<&NamespaceIdent>,
157 ) -> iceberg::Result<Vec<NamespaceIdent>> {
158 let inner = self.inner.clone();
159 execute_blocking_jni(move || {
160 execute_with_jni_env(inner.jvm, |env| {
161 let result_json =
162 call_method!(env, inner.java_catalog.as_obj(), {String listNamespaces()})
163 .with_context(|| "Failed to list iceberg namespaces".to_owned())?;
164
165 let rust_json_str = jobj_to_str(env, result_json)?;
166
167 let resp: ListNamespacesResponse = serde_json::from_str(&rust_json_str)?;
168
169 Ok(resp.namespaces)
170 })
171 })
172 .await
173 .map_err(|e| {
174 iceberg::Error::new(
175 iceberg::ErrorKind::Unexpected,
176 "Failed to list iceberg namespaces.",
177 )
178 .with_source(e)
179 })
180 }
181
182 async fn create_namespace(
184 &self,
185 namespace: &iceberg::NamespaceIdent,
186 _properties: HashMap<String, String>,
187 ) -> iceberg::Result<iceberg::Namespace> {
188 let inner = self.inner.clone();
189 let namespace = namespace.clone();
190 execute_blocking_jni(move || {
191 execute_with_jni_env(inner.jvm, |env| {
192 let namespace_str = namespace_to_string(&namespace);
193 let namespace_jstr = env.new_string(&namespace_str).unwrap();
194
195 call_method!(env, inner.java_catalog.as_obj(), {void createNamespace(String)},
196 &namespace_jstr)
197 .with_context(|| format!("Failed to create namespace: {namespace}"))?;
198
199 Ok(Namespace::new(namespace))
200 })
201 })
202 .await
203 .map_err(|e| {
204 iceberg::Error::new(
205 iceberg::ErrorKind::Unexpected,
206 "Failed to create namespace.",
207 )
208 .with_source(e)
209 })
210 }
211
212 async fn get_namespace(&self, namespace: &NamespaceIdent) -> iceberg::Result<Namespace> {
214 let inner = self.inner.clone();
215 let namespace = namespace.clone();
216 execute_blocking_jni(move || {
217 execute_with_jni_env(inner.jvm, |env| {
218 let namespace_jstr = env.new_string(namespace_to_string(&namespace))?;
219 let result_json =
220 call_method!(env, inner.java_catalog.as_obj(), {String getNamespace(String)},
221 &namespace_jstr)
222 .with_context(|| format!("Failed to get iceberg namespace: {namespace}"))?;
223 let rust_json_str = jobj_to_str(env, result_json)?;
224 let resp: GetNamespaceResponse = serde_json::from_str(&rust_json_str)?;
225 Ok(Namespace::with_properties(resp.namespace, resp.properties))
226 })
227 })
228 .await
229 .map_err(|e| {
230 iceberg::Error::new(
231 iceberg::ErrorKind::Unexpected,
232 "Failed to get iceberg namespace.",
233 )
234 .with_source(e)
235 })
236 }
237
238 async fn namespace_exists(&self, namespace: &NamespaceIdent) -> iceberg::Result<bool> {
240 let inner = self.inner.clone();
241 let namespace = namespace.clone();
242 execute_blocking_jni(move || {
243 execute_with_jni_env(inner.jvm, |env| {
244 let namespace_str = namespace_to_string(&namespace);
245 let namespace_jstr = env.new_string(&namespace_str).unwrap();
246
247 let exists =
248 call_method!(env, inner.java_catalog.as_obj(), {boolean namespaceExists(String)},
249 &namespace_jstr)
250 .with_context(|| format!("Failed to check namespace exists: {namespace}"))?;
251
252 Ok(exists)
253 })
254 })
255 .await
256 .map_err(|e| {
257 iceberg::Error::new(
258 iceberg::ErrorKind::Unexpected,
259 "Failed to check namespace exists.",
260 )
261 .with_source(e)
262 })
263 }
264
265 async fn drop_namespace(&self, _namespace: &NamespaceIdent) -> iceberg::Result<()> {
267 todo!()
268 }
269
270 async fn list_tables(&self, namespace: &NamespaceIdent) -> iceberg::Result<Vec<TableIdent>> {
272 let inner = self.inner.clone();
273 let namespace = namespace.clone();
274 execute_blocking_jni(move || {
275 execute_with_jni_env(inner.jvm, |env| {
276 let namespace_str = namespace_to_string(&namespace);
277 let namespace_jstr = env.new_string(&namespace_str).unwrap();
278
279 let result_json =
280 call_method!(env, inner.java_catalog.as_obj(), {String listTables(String)},
281 &namespace_jstr)
282 .with_context(|| {
283 format!("Failed to list iceberg tables in namespace: {}", namespace)
284 })?;
285
286 let rust_json_str = jobj_to_str(env, result_json)?;
287
288 let resp: ListTablesResponse = serde_json::from_str(&rust_json_str)?;
289
290 Ok(resp.identifiers)
291 })
292 })
293 .await
294 .map_err(|e| {
295 iceberg::Error::new(
296 iceberg::ErrorKind::Unexpected,
297 "Failed to list iceberg tables.",
298 )
299 .with_source(e)
300 })
301 }
302
303 async fn update_namespace(
304 &self,
305 _namespace: &NamespaceIdent,
306 _properties: HashMap<String, String>,
307 ) -> iceberg::Result<()> {
308 todo!()
309 }
310
311 async fn create_table(
313 &self,
314 namespace: &NamespaceIdent,
315 creation: TableCreation,
316 ) -> iceberg::Result<Table> {
317 let inner = self.inner.clone();
318 let file_io_props = self.file_io_props.clone();
319 let runtime = Runtime::try_current()?;
320 let namespace = namespace.clone();
321 execute_blocking_jni(move || {
322 execute_with_jni_env(inner.jvm, |env| {
323 let namespace_str = namespace_to_string(&namespace);
324 let namespace_jstr = env.new_string(&namespace_str).unwrap();
325
326 let creation_str = serde_json::to_string(&CreateTableRequest::from(&creation))?;
327
328 let creation_jstr = env.new_string(&creation_str).unwrap();
329
330 let result_json =
331 call_method!(env, inner.java_catalog.as_obj(), {String createTable(String, String)},
332 &namespace_jstr, &creation_jstr)
333 .with_context(|| {
334 format!("Failed to create iceberg table: {}", creation.name)
335 })?;
336
337 let rust_json_str = jobj_to_str(env, result_json)?;
338
339 let resp: LoadTableResponse = serde_json::from_str(&rust_json_str)?;
340
341 let _metadata_location = resp.metadata_location.ok_or_else(|| {
342 iceberg::Error::new(
343 iceberg::ErrorKind::FeatureUnsupported,
344 "Loading uncommitted table is not supported!",
345 )
346 })?;
347
348 let table_metadata = resp.metadata;
349
350 let file_io =
351 FileIOBuilder::new(Arc::new(OpenDalResolvingStorageFactory::new()))
352 .with_props(file_io_props.iter())
353 .build();
354
355 Ok(Table::builder()
356 .file_io(file_io)
357 .identifier(TableIdent::new(namespace, creation.name))
358 .metadata(table_metadata)
359 .runtime(runtime)
360 .build())
361 })
362 })
363 .await
364 .map_err(|e| {
365 iceberg::Error::new(
366 iceberg::ErrorKind::Unexpected,
367 "Failed to create iceberg table.",
368 )
369 .with_source(e)
370 })?
371 }
372
373 async fn load_table(&self, table: &TableIdent) -> iceberg::Result<Table> {
375 let inner = self.inner.clone();
376 let file_io_props = self.file_io_props.clone();
377 let runtime = Runtime::try_current()?;
378 let table = table.clone();
379 execute_blocking_jni(move || {
380 execute_with_jni_env(inner.jvm, |env| {
381 let table_name_str = table.to_string();
382
383 let table_name_jstr = env.new_string(&table_name_str).unwrap();
384
385 let result_json =
386 call_method!(env, inner.java_catalog.as_obj(), {String loadTable(String)},
387 &table_name_jstr)
388 .with_context(|| format!("Failed to load iceberg table: {table_name_str}"))?;
389
390 let rust_json_str = jobj_to_str(env, result_json)?;
391
392 let resp: LoadTableResponse = serde_json::from_str(&rust_json_str)?;
393
394 let metadata_location = resp.metadata_location.ok_or_else(|| {
395 iceberg::Error::new(
396 iceberg::ErrorKind::FeatureUnsupported,
397 "Loading uncommitted table is not supported!",
398 )
399 })?;
400
401 tracing::info!(
402 "Table metadata location of {table_name_str} is {metadata_location}"
403 );
404
405 let table_metadata = resp.metadata;
406
407 let file_io = FileIOBuilder::new(Arc::new(OpenDalResolvingStorageFactory::new()))
408 .with_props(file_io_props.iter())
409 .build();
410
411 Ok(Table::builder()
412 .file_io(file_io)
413 .identifier(table)
414 .metadata(table_metadata)
415 .runtime(runtime)
416 .build())
417 })
418 })
419 .await
420 .map_err(|e| {
421 iceberg::Error::new(
422 iceberg::ErrorKind::Unexpected,
423 "Failed to load iceberg table.",
424 )
425 .with_source(e)
426 })?
427 }
428
429 async fn drop_table(&self, table: &TableIdent) -> iceberg::Result<()> {
431 let inner = self.inner.clone();
432 let table = table.to_owned();
433 execute_blocking_jni(move || {
434 execute_with_jni_env(inner.jvm, |env| {
435 let table_name_str = table.to_string();
436
437 let table_name_jstr = env.new_string(&table_name_str).unwrap();
438
439 call_method!(env, inner.java_catalog.as_obj(), {boolean dropTable(String)},
440 &table_name_jstr)
441 .with_context(|| format!("Failed to drop iceberg table: {table_name_str}"))?;
442
443 Ok(())
444 })
445 })
446 .await
447 .map_err(|e| {
448 iceberg::Error::new(
449 iceberg::ErrorKind::Unexpected,
450 "Failed to drop iceberg table.",
451 )
452 .with_source(e)
453 })
454 }
455
456 async fn purge_table(&self, table: &TableIdent) -> iceberg::Result<()> {
457 let table_info = self.load_table(table).await?;
458 self.drop_table(table).await?;
459 iceberg::drop_table_data(&table_info).await
460 }
461
462 async fn register_table(
463 &self,
464 _table_ident: &TableIdent,
465 _metadata_location: String,
466 ) -> iceberg::Result<Table> {
467 Err(iceberg::Error::new(
468 iceberg::ErrorKind::Unexpected,
469 "register_table is not supported by JniCatalog",
470 ))
471 }
472
473 async fn table_exists(&self, table: &TableIdent) -> iceberg::Result<bool> {
475 let inner = self.inner.clone();
476 let table = table.clone();
477 execute_blocking_jni(move || {
478 execute_with_jni_env(inner.jvm, |env| {
479 let table_name_str = table.to_string();
480
481 let table_name_jstr = env.new_string(&table_name_str).unwrap();
482
483 let exists =
484 call_method!(env, inner.java_catalog.as_obj(), {boolean tableExists(String)},
485 &table_name_jstr)
486 .with_context(|| {
487 format!("Failed to check iceberg table exists: {table_name_str}")
488 })?;
489
490 Ok(exists)
491 })
492 })
493 .await
494 .map_err(|e| {
495 iceberg::Error::new(
496 iceberg::ErrorKind::Unexpected,
497 "Failed to check iceberg table exists.",
498 )
499 .with_source(e)
500 })
501 }
502
503 async fn rename_table(&self, _src: &TableIdent, _dest: &TableIdent) -> iceberg::Result<()> {
505 todo!()
506 }
507
508 async fn update_table(&self, mut commit: TableCommit) -> iceberg::Result<Table> {
510 let inner = self.inner.clone();
511 let file_io_props = self.file_io_props.clone();
512 let runtime = Runtime::try_current()?;
513 execute_blocking_jni(move || {
514 execute_with_jni_env(inner.jvm, |env| {
515 let requirements = commit.take_requirements();
516 let updates = commit.take_updates();
517 let request = CommitTableRequest {
518 identifier: commit.identifier().clone(),
519 requirements,
520 updates,
521 };
522 let request_str = serde_json::to_string(&request)?;
523
524 let request_jni_str = env.new_string(&request_str).with_context(|| {
525 format!("Failed to create jni string from request json: {request_str}.")
526 })?;
527
528 let result_json =
529 call_method!(env, inner.java_catalog.as_obj(), {String updateTable(String)},
530 &request_jni_str)
531 .with_context(|| {
532 format!("Failed to update iceberg table: {}", commit.identifier())
533 })?;
534
535 let rust_json_str = jobj_to_str(env, result_json)?;
536
537 let response: CommitTableResponse = serde_json::from_str(&rust_json_str)?;
538
539 tracing::info!(
540 "Table metadata location of {} is {}",
541 commit.identifier(),
542 response.metadata_location
543 );
544
545 let table_metadata = response.metadata;
546
547 let file_io = FileIOBuilder::new(Arc::new(OpenDalResolvingStorageFactory::new()))
548 .with_props(file_io_props.iter())
549 .build();
550
551 Ok(Table::builder()
552 .file_io(file_io)
553 .identifier(commit.identifier().clone())
554 .metadata(table_metadata)
555 .runtime(runtime)
556 .build()?)
557 })
558 })
559 .await
560 .map_err(|e| {
561 iceberg::Error::new(
562 iceberg::ErrorKind::Unexpected,
563 "Failed to update iceberg table.",
564 )
565 .with_source(e)
566 })
567 }
568}
569
570impl Drop for JniCatalogInner {
571 fn drop(&mut self) {
572 let _ = execute_with_jni_env(self.jvm, |env| {
573 call_method!(env, self.java_catalog.as_obj(), {void close()})
574 .with_context(|| "Failed to close iceberg catalog".to_owned())?;
575 Ok(())
576 })
577 .inspect_err(
578 |e| tracing::error!(error = ?e.as_report(), "Failed to close iceberg catalog"),
579 );
580 }
581}
582
583impl JniCatalog {
584 fn build(
585 file_io_props: HashMap<String, String>,
586 name: impl ToString,
587 catalog_impl: impl ToString,
588 java_catalog_props: HashMap<String, String>,
589 ) -> ConnectorResult<Self> {
590 let jvm = Jvm::get_or_init()?;
591
592 execute_with_jni_env(jvm, |env| {
593 let props = env.new_object_array(
595 (java_catalog_props.len() * 2) as i32,
596 "java/lang/String",
597 JObject::null(),
598 )?;
599 for (i, (key, value)) in java_catalog_props.iter().enumerate() {
600 let key_j_str = env.new_string(key)?;
601 let value_j_str = env.new_string(value)?;
602 env.set_object_array_element(&props, i as i32 * 2, key_j_str)?;
603 env.set_object_array_element(&props, i as i32 * 2 + 1, value_j_str)?;
604 }
605
606 let jni_catalog_wrapper = env
607 .call_static_method(
608 "com/risingwave/connector/catalog/JniCatalogWrapper",
609 "create",
610 "(Ljava/lang/String;Ljava/lang/String;[Ljava/lang/String;)Lcom/risingwave/connector/catalog/JniCatalogWrapper;",
611 &[
612 (&env.new_string(name.to_string()).unwrap()).into(),
613 (&env.new_string(catalog_impl.to_string()).unwrap()).into(),
614 (&props).into(),
615 ],
616 )?;
617
618 let jni_catalog = env.new_global_ref(jni_catalog_wrapper.l().unwrap())?;
619
620 Ok(Self {
621 inner: Arc::new(JniCatalogInner {
622 java_catalog: jni_catalog,
623 jvm,
624 }),
625 file_io_props: Arc::new(file_io_props),
626 })
627 })
628 .map_err(Into::into)
629 }
630
631 pub async fn build_catalog(
632 file_io_props: HashMap<String, String>,
633 name: impl ToString + Send + 'static,
634 catalog_impl: impl ToString + Send + 'static,
635 java_catalog_props: HashMap<String, String>,
636 ) -> ConnectorResult<Arc<dyn Catalog>> {
637 let catalog = execute_blocking_jni(move || {
638 Ok(Self::build(
639 file_io_props,
640 name,
641 catalog_impl,
642 java_catalog_props,
643 )?)
644 })
645 .await?;
646 Ok(Arc::new(catalog) as Arc<dyn Catalog>)
647 }
648}