Skip to main content

risingwave_connector/connector_common/iceberg/
jni_catalog.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! This module provide jni catalog.
16
17#![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    /// The name of the table.
58    pub name: String,
59    /// The location of the table.
60    pub location: Option<String>,
61    /// The schema of the table.
62    pub schema: Schema,
63    /// The partition spec of the table, could be None.
64    pub partition_spec: Option<UnboundPartitionSpec>,
65    /// The sort order of the table.
66    pub write_order: Option<SortOrder>,
67    /// The properties of the table.
68    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    /// A blocking JNI operation can outlive its async caller after cancellation. Keep the Java
132    /// catalog alive until the last such operation finishes.
133    inner: Arc<JniCatalogInner>,
134    file_io_props: Arc<HashMap<String, String>>,
135}
136
137/// Iceberg's Java catalog API is synchronous and may perform remote catalog I/O. Running it
138/// directly in an async `Catalog` method would block a Tokio worker thread, so async JNI operations
139/// go through the runtime's blocking pool.
140async 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    /// List namespaces from the catalog.
154    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    /// Create a new namespace inside the catalog.
183    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    /// Get a namespace information from the catalog.
213    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    /// Check if namespace exists in catalog.
239    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    /// Drop a namespace from the catalog.
266    async fn drop_namespace(&self, _namespace: &NamespaceIdent) -> iceberg::Result<()> {
267        todo!()
268    }
269
270    /// List tables from namespace.
271    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    /// Create a new table inside the namespace.
312    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    /// Load table from the catalog.
374    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    /// Drop a table from the catalog.
430    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    /// Check if a table exists in the catalog.
474    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    /// Rename a table in the catalog.
504    async fn rename_table(&self, _src: &TableIdent, _dest: &TableIdent) -> iceberg::Result<()> {
505        todo!()
506    }
507
508    /// Update a table to the catalog.
509    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            // Convert props to string array
594            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}