public abstract class OceanBaseCatalog extends Object implements Serializable
OceanBaseCatalog for OceanBase connector that supports schema evolution.| 限定符和类型 | 字段和说明 |
|---|---|
protected com.oceanbase.connector.flink.connection.OceanBaseConnectionProvider |
connectionProvider |
| 构造器和说明 |
|---|
OceanBaseCatalog(com.oceanbase.connector.flink.OceanBaseConnectorOptions connectorOptions) |
| 限定符和类型 | 方法和说明 |
|---|---|
abstract void |
alterAddColumns(String databaseName,
String tableName,
List<OceanBaseColumn> addColumns) |
abstract void |
alterColumnType(String schemaName,
String tableName,
String columnName,
org.apache.flink.cdc.common.types.DataType dataType) |
abstract void |
alterDropColumns(String schemaName,
String tableName,
List<String> dropColumns) |
void |
close() |
abstract void |
createDatabase(String databaseName,
boolean ignoreIfExists) |
abstract void |
createTable(OceanBaseTable table,
boolean ignoreIfExists) |
abstract boolean |
databaseExists(String databaseName) |
abstract void |
dropTable(String schemaName,
String tableName) |
protected List<String> |
executeSingleColumnStatement(String sql) |
protected void |
executeUpdateStatement(String sql) |
void |
open() |
abstract void |
renameColumn(String schemaName,
String tableName,
String oldColumnName,
String newColumnName) |
abstract boolean |
tableExists(String databaseName,
String tableName) |
abstract void |
truncateTable(String schemaName,
String tableName) |
protected com.oceanbase.connector.flink.connection.OceanBaseConnectionProvider connectionProvider
public OceanBaseCatalog(com.oceanbase.connector.flink.OceanBaseConnectorOptions connectorOptions)
public void open()
protected List<String> executeSingleColumnStatement(String sql) throws SQLException
SQLExceptionprotected void executeUpdateStatement(String sql) throws SQLException
SQLExceptionpublic abstract boolean databaseExists(String databaseName) throws OceanBaseCatalogException
public abstract void createDatabase(String databaseName, boolean ignoreIfExists) throws OceanBaseCatalogException
public abstract boolean tableExists(String databaseName, String tableName) throws OceanBaseCatalogException
public abstract void createTable(OceanBaseTable table, boolean ignoreIfExists) throws OceanBaseCatalogException
public abstract void alterAddColumns(String databaseName, String tableName, List<OceanBaseColumn> addColumns)
public abstract void alterDropColumns(String schemaName, String tableName, List<String> dropColumns)
public abstract void alterColumnType(String schemaName, String tableName, String columnName, org.apache.flink.cdc.common.types.DataType dataType)
public abstract void renameColumn(String schemaName, String tableName, String oldColumnName, String newColumnName)
public void close()
Copyright © 2025 The Apache Software Foundation. All rights reserved.