Documentation
¶
Index ¶
- func NewAdapter() adapter.DatabaseAdapter
- type Adapter
- func (a *Adapter) Capabilities() dbcapabilities.Capability
- func (a *Adapter) Connect(ctx context.Context, config adapter.ConnectionConfig) (adapter.Connection, error)
- func (a *Adapter) ConnectInstance(ctx context.Context, config adapter.InstanceConfig) (adapter.InstanceConnection, error)
- func (a *Adapter) Type() dbcapabilities.DatabaseType
- type Connection
- func (c *Connection) Adapter() adapter.DatabaseAdapter
- func (c *Connection) Close() error
- func (c *Connection) Config() adapter.ConnectionConfig
- func (c *Connection) DataOperations() adapter.DataOperator
- func (c *Connection) ID() string
- func (c *Connection) IsConnected() bool
- func (c *Connection) MetadataOperations() adapter.MetadataOperator
- func (c *Connection) Ping(ctx context.Context) error
- func (c *Connection) Raw() interface{}
- func (c *Connection) ReplicationOperations() adapter.ReplicationOperator
- func (c *Connection) SchemaOperations() adapter.SchemaOperator
- func (c *Connection) Type() dbcapabilities.DatabaseType
- type DataOps
- func (d *DataOps) Delete(ctx context.Context, table string, conditions map[string]interface{}) (int64, error)
- func (d *DataOps) ExecuteCountQuery(ctx context.Context, query string) (int64, error)
- func (d *DataOps) ExecuteQuery(ctx context.Context, query string, args ...interface{}) ([]interface{}, error)
- func (d *DataOps) Fetch(ctx context.Context, table string, limit int) ([]map[string]interface{}, error)
- func (d *DataOps) FetchWithColumns(ctx context.Context, table string, columns []string, limit int) ([]map[string]interface{}, error)
- func (d *DataOps) GetRowCount(ctx context.Context, table string, whereClause string) (int64, bool, error)
- func (d *DataOps) Insert(ctx context.Context, table string, data []map[string]interface{}) (int64, error)
- func (d *DataOps) Stream(ctx context.Context, params adapter.StreamParams) (adapter.StreamResult, error)
- func (d *DataOps) Update(ctx context.Context, table string, data []map[string]interface{}, ...) (int64, error)
- func (d *DataOps) Upsert(ctx context.Context, table string, data []map[string]interface{}, ...) (int64, error)
- func (d *DataOps) Wipe(ctx context.Context) error
- type DatabricksClient
- func (c *DatabricksClient) Close() error
- func (c *DatabricksClient) CreateDatabase(ctx context.Context, name string, options map[string]interface{}) error
- func (c *DatabricksClient) DB() *sql.DB
- func (c *DatabricksClient) DropDatabase(ctx context.Context, name string) error
- func (c *DatabricksClient) GetDatabaseName() string
- func (c *DatabricksClient) ListDatabases(ctx context.Context) ([]string, error)
- func (c *DatabricksClient) Ping(ctx context.Context) error
- type InstanceConnection
- func (ic *InstanceConnection) Adapter() adapter.DatabaseAdapter
- func (ic *InstanceConnection) Close() error
- func (ic *InstanceConnection) Config() adapter.InstanceConfig
- func (ic *InstanceConnection) CreateDatabase(ctx context.Context, name string, options map[string]interface{}) error
- func (ic *InstanceConnection) DropDatabase(ctx context.Context, name string, options map[string]interface{}) error
- func (ic *InstanceConnection) ID() string
- func (ic *InstanceConnection) IsConnected() bool
- func (ic *InstanceConnection) ListDatabases(ctx context.Context) ([]string, error)
- func (ic *InstanceConnection) MetadataOperations() adapter.MetadataOperator
- func (ic *InstanceConnection) Ping(ctx context.Context) error
- func (ic *InstanceConnection) Raw() interface{}
- func (ic *InstanceConnection) Type() dbcapabilities.DatabaseType
- type MetadataOps
- func (m *MetadataOps) CollectDatabaseMetadata(ctx context.Context) (map[string]interface{}, error)
- func (m *MetadataOps) CollectInstanceMetadata(ctx context.Context) (map[string]interface{}, error)
- func (m *MetadataOps) CollectInstanceMetrics(ctx context.Context) (map[string]interface{}, error)
- func (m *MetadataOps) ExecuteCommand(ctx context.Context, command string) ([]byte, error)
- func (m *MetadataOps) GetDatabaseSize(ctx context.Context) (int64, error)
- func (m *MetadataOps) GetTableCount(ctx context.Context) (int, error)
- func (m *MetadataOps) GetUniqueIdentifier(ctx context.Context) (string, error)
- func (m *MetadataOps) GetVersion(ctx context.Context) (string, error)
- func (m *MetadataOps) ListLogicalDatabases(ctx context.Context) ([]adapter.LogicalDatabaseInfo, error)
- type ReplicationOps
- func (r *ReplicationOps) ApplyCDCEvent(ctx context.Context, event *adapter.CDCEvent) error
- func (r *ReplicationOps) CheckPrerequisites(ctx context.Context) error
- func (r *ReplicationOps) Connect(ctx context.Context, config adapter.ReplicationConfig) (adapter.ReplicationSource, error)
- func (r *ReplicationOps) DropPublication(ctx context.Context, publicationName string) error
- func (r *ReplicationOps) DropSlot(ctx context.Context, slotName string) error
- func (r *ReplicationOps) GetLag(ctx context.Context) (map[string]interface{}, error)
- func (r *ReplicationOps) GetStatus(ctx context.Context) (map[string]interface{}, error)
- func (r *ReplicationOps) GetSupportedMechanisms() []string
- func (r *ReplicationOps) IsSupported() bool
- func (r *ReplicationOps) ListPublications(ctx context.Context) ([]map[string]interface{}, error)
- func (r *ReplicationOps) ListSlots(ctx context.Context) ([]map[string]interface{}, error)
- func (r *ReplicationOps) ParseEvent(ctx context.Context, rawEvent map[string]interface{}) (*adapter.CDCEvent, error)
- func (r *ReplicationOps) TransformData(ctx context.Context, data map[string]interface{}, ...) (map[string]interface{}, error)
- type SchemaOps
- func (s *SchemaOps) CreateStructure(ctx context.Context, model *unifiedmodel.UnifiedModel) error
- func (s *SchemaOps) DiscoverSchema(ctx context.Context) (*unifiedmodel.UnifiedModel, error)
- func (s *SchemaOps) GetTableSchema(ctx context.Context, tableName string) (*unifiedmodel.Table, error)
- func (s *SchemaOps) ListTables(ctx context.Context) ([]string, error)
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func NewAdapter ¶
func NewAdapter() adapter.DatabaseAdapter
NewAdapter creates a new Databricks adapter instance.
Types ¶
type Adapter ¶
type Adapter struct{}
Adapter implements adapter.DatabaseAdapter for Databricks.
func (*Adapter) Capabilities ¶
func (a *Adapter) Capabilities() dbcapabilities.Capability
Capabilities returns the capability metadata.
func (*Adapter) Connect ¶
func (a *Adapter) Connect(ctx context.Context, config adapter.ConnectionConfig) (adapter.Connection, error)
Connect establishes a connection to Databricks.
func (*Adapter) ConnectInstance ¶
func (a *Adapter) ConnectInstance(ctx context.Context, config adapter.InstanceConfig) (adapter.InstanceConnection, error)
ConnectInstance establishes an instance-level connection to Databricks.
func (*Adapter) Type ¶
func (a *Adapter) Type() dbcapabilities.DatabaseType
Type returns the database type identifier.
type Connection ¶
type Connection struct {
// contains filtered or unexported fields
}
Connection implements adapter.Connection for Databricks.
func (*Connection) Adapter ¶
func (c *Connection) Adapter() adapter.DatabaseAdapter
Adapter returns the database adapter.
func (*Connection) Config ¶
func (c *Connection) Config() adapter.ConnectionConfig
Config returns the connection configuration.
func (*Connection) DataOperations ¶
func (c *Connection) DataOperations() adapter.DataOperator
DataOperations returns the data operator.
func (*Connection) IsConnected ¶
func (c *Connection) IsConnected() bool
IsConnected returns whether the connection is active.
func (*Connection) MetadataOperations ¶
func (c *Connection) MetadataOperations() adapter.MetadataOperator
MetadataOperations returns the metadata operator.
func (*Connection) Ping ¶
func (c *Connection) Ping(ctx context.Context) error
Ping tests the connection.
func (*Connection) Raw ¶
func (c *Connection) Raw() interface{}
Raw returns the underlying Databricks client.
func (*Connection) ReplicationOperations ¶
func (c *Connection) ReplicationOperations() adapter.ReplicationOperator
ReplicationOperations returns the replication operator.
func (*Connection) SchemaOperations ¶
func (c *Connection) SchemaOperations() adapter.SchemaOperator
SchemaOperations returns the schema operator.
func (*Connection) Type ¶
func (c *Connection) Type() dbcapabilities.DatabaseType
Type returns the database type.
type DataOps ¶
type DataOps struct {
// contains filtered or unexported fields
}
DataOps implements data operations for Databricks.
func (*DataOps) Delete ¶
func (d *DataOps) Delete(ctx context.Context, table string, conditions map[string]interface{}) (int64, error)
Delete deletes rows from Databricks.
func (*DataOps) ExecuteCountQuery ¶
ExecuteCountQuery counts rows.
func (*DataOps) ExecuteQuery ¶
func (d *DataOps) ExecuteQuery(ctx context.Context, query string, args ...interface{}) ([]interface{}, error)
ExecuteQuery executes a SQL query.
func (*DataOps) Fetch ¶
func (d *DataOps) Fetch(ctx context.Context, table string, limit int) ([]map[string]interface{}, error)
Fetch retrieves rows from Databricks.
func (*DataOps) FetchWithColumns ¶
func (d *DataOps) FetchWithColumns(ctx context.Context, table string, columns []string, limit int) ([]map[string]interface{}, error)
FetchWithColumns retrieves rows with specific columns.
func (*DataOps) GetRowCount ¶
func (d *DataOps) GetRowCount(ctx context.Context, table string, whereClause string) (int64, bool, error)
GetRowCount returns the number of rows in a table.
func (*DataOps) Insert ¶
func (d *DataOps) Insert(ctx context.Context, table string, data []map[string]interface{}) (int64, error)
Insert inserts rows into Databricks.
func (*DataOps) Stream ¶
func (d *DataOps) Stream(ctx context.Context, params adapter.StreamParams) (adapter.StreamResult, error)
Stream retrieves rows in batches.
func (*DataOps) Update ¶
func (d *DataOps) Update(ctx context.Context, table string, data []map[string]interface{}, whereColumns []string) (int64, error)
Update updates rows in Databricks.
type DatabricksClient ¶
type DatabricksClient struct {
// contains filtered or unexported fields
}
DatabricksClient wraps the SQL database connection for Databricks.
func NewDatabricksClient ¶
func NewDatabricksClient(ctx context.Context, cfg adapter.ConnectionConfig) (*DatabricksClient, error)
NewDatabricksClient creates a new Databricks client from a database connection config.
func NewDatabricksClientFromInstance ¶
func NewDatabricksClientFromInstance(ctx context.Context, cfg adapter.InstanceConfig) (*DatabricksClient, error)
NewDatabricksClientFromInstance creates a new Databricks client from an instance config.
func (*DatabricksClient) Close ¶
func (c *DatabricksClient) Close() error
Close closes the Databricks client.
func (*DatabricksClient) CreateDatabase ¶
func (c *DatabricksClient) CreateDatabase(ctx context.Context, name string, options map[string]interface{}) error
CreateDatabase creates a new database (schema) in Databricks.
func (*DatabricksClient) DB ¶
func (c *DatabricksClient) DB() *sql.DB
DB returns the underlying database connection.
func (*DatabricksClient) DropDatabase ¶
func (c *DatabricksClient) DropDatabase(ctx context.Context, name string) error
DropDatabase drops a database (schema) from Databricks.
func (*DatabricksClient) GetDatabaseName ¶
func (c *DatabricksClient) GetDatabaseName() string
GetDatabaseName returns the current database name.
func (*DatabricksClient) ListDatabases ¶
func (c *DatabricksClient) ListDatabases(ctx context.Context) ([]string, error)
ListDatabases lists all databases (schemas) in Databricks.
type InstanceConnection ¶
type InstanceConnection struct {
// contains filtered or unexported fields
}
InstanceConnection implements adapter.InstanceConnection for Databricks.
func (*InstanceConnection) Adapter ¶
func (ic *InstanceConnection) Adapter() adapter.DatabaseAdapter
Adapter returns the database adapter.
func (*InstanceConnection) Close ¶
func (ic *InstanceConnection) Close() error
Close closes the connection.
func (*InstanceConnection) Config ¶
func (ic *InstanceConnection) Config() adapter.InstanceConfig
Config returns the instance configuration.
func (*InstanceConnection) CreateDatabase ¶
func (ic *InstanceConnection) CreateDatabase(ctx context.Context, name string, options map[string]interface{}) error
CreateDatabase creates a new Databricks database (schema).
func (*InstanceConnection) DropDatabase ¶
func (ic *InstanceConnection) DropDatabase(ctx context.Context, name string, options map[string]interface{}) error
DropDatabase deletes a Databricks database (schema).
func (*InstanceConnection) ID ¶
func (ic *InstanceConnection) ID() string
ID returns the instance connection identifier.
func (*InstanceConnection) IsConnected ¶
func (ic *InstanceConnection) IsConnected() bool
IsConnected returns whether the connection is active.
func (*InstanceConnection) ListDatabases ¶
func (ic *InstanceConnection) ListDatabases(ctx context.Context) ([]string, error)
ListDatabases lists all Databricks databases (schemas).
func (*InstanceConnection) MetadataOperations ¶
func (ic *InstanceConnection) MetadataOperations() adapter.MetadataOperator
MetadataOperations returns the metadata operator.
func (*InstanceConnection) Ping ¶
func (ic *InstanceConnection) Ping(ctx context.Context) error
Ping tests the connection.
func (*InstanceConnection) Raw ¶
func (ic *InstanceConnection) Raw() interface{}
Raw returns the underlying Databricks client.
func (*InstanceConnection) Type ¶
func (ic *InstanceConnection) Type() dbcapabilities.DatabaseType
Type returns the database type.
type MetadataOps ¶
type MetadataOps struct {
// contains filtered or unexported fields
}
MetadataOps implements metadata operations for Databricks.
func (*MetadataOps) CollectDatabaseMetadata ¶
func (m *MetadataOps) CollectDatabaseMetadata(ctx context.Context) (map[string]interface{}, error)
CollectDatabaseMetadata collects metadata about the Databricks database.
func (*MetadataOps) CollectInstanceMetadata ¶
func (m *MetadataOps) CollectInstanceMetadata(ctx context.Context) (map[string]interface{}, error)
CollectInstanceMetadata collects metadata about the Databricks workspace.
func (*MetadataOps) CollectInstanceMetrics ¶
func (m *MetadataOps) CollectInstanceMetrics(ctx context.Context) (map[string]interface{}, error)
CollectInstanceMetrics collects performance metrics.
func (*MetadataOps) ExecuteCommand ¶
ExecuteCommand executes a Databricks SQL command.
func (*MetadataOps) GetDatabaseSize ¶
func (m *MetadataOps) GetDatabaseSize(ctx context.Context) (int64, error)
GetDatabaseSize returns the size of the database in bytes.
func (*MetadataOps) GetTableCount ¶
func (m *MetadataOps) GetTableCount(ctx context.Context) (int, error)
GetTableCount returns the number of tables in the database.
func (*MetadataOps) GetUniqueIdentifier ¶
func (m *MetadataOps) GetUniqueIdentifier(ctx context.Context) (string, error)
GetUniqueIdentifier returns the workspace identifier.
func (*MetadataOps) GetVersion ¶
func (m *MetadataOps) GetVersion(ctx context.Context) (string, error)
GetVersion returns the Databricks runtime version.
func (*MetadataOps) ListLogicalDatabases ¶
func (m *MetadataOps) ListLogicalDatabases(ctx context.Context) ([]adapter.LogicalDatabaseInfo, error)
ListLogicalDatabases lists logical databases.
type ReplicationOps ¶
type ReplicationOps struct {
// contains filtered or unexported fields
}
ReplicationOps implements replication operations for Databricks.
func (*ReplicationOps) ApplyCDCEvent ¶
ApplyCDCEvent applies a CDC event (not implemented).
func (*ReplicationOps) CheckPrerequisites ¶
func (r *ReplicationOps) CheckPrerequisites(ctx context.Context) error
CheckPrerequisites checks if prerequisites for CDC are met.
func (*ReplicationOps) Connect ¶
func (r *ReplicationOps) Connect(ctx context.Context, config adapter.ReplicationConfig) (adapter.ReplicationSource, error)
Connect establishes a CDC connection.
func (*ReplicationOps) DropPublication ¶
func (r *ReplicationOps) DropPublication(ctx context.Context, publicationName string) error
DropPublication drops a publication (not applicable for Databricks).
func (*ReplicationOps) DropSlot ¶
func (r *ReplicationOps) DropSlot(ctx context.Context, slotName string) error
DropSlot drops a replication slot (not applicable for Databricks).
func (*ReplicationOps) GetLag ¶
func (r *ReplicationOps) GetLag(ctx context.Context) (map[string]interface{}, error)
GetLag returns the replication lag.
func (*ReplicationOps) GetStatus ¶
func (r *ReplicationOps) GetStatus(ctx context.Context) (map[string]interface{}, error)
GetStatus returns the CDC status.
func (*ReplicationOps) GetSupportedMechanisms ¶
func (r *ReplicationOps) GetSupportedMechanisms() []string
GetSupportedMechanisms returns the list of supported CDC mechanisms.
func (*ReplicationOps) IsSupported ¶
func (r *ReplicationOps) IsSupported() bool
IsSupported returns whether CDC/replication is supported.
func (*ReplicationOps) ListPublications ¶
func (r *ReplicationOps) ListPublications(ctx context.Context) ([]map[string]interface{}, error)
ListPublications lists publications (not applicable for Databricks).
func (*ReplicationOps) ListSlots ¶
func (r *ReplicationOps) ListSlots(ctx context.Context) ([]map[string]interface{}, error)
ListSlots lists replication slots (not applicable for Databricks).
func (*ReplicationOps) ParseEvent ¶
func (r *ReplicationOps) ParseEvent(ctx context.Context, rawEvent map[string]interface{}) (*adapter.CDCEvent, error)
ParseEvent parses a CDC event (not implemented).
func (*ReplicationOps) TransformData ¶
func (r *ReplicationOps) TransformData(ctx context.Context, data map[string]interface{}, rules []adapter.TransformationRule, transformationServiceEndpoint string) (map[string]interface{}, error)
TransformData transforms data using transformation rules.
type SchemaOps ¶
type SchemaOps struct {
// contains filtered or unexported fields
}
SchemaOps implements schema operations for Databricks.
func (*SchemaOps) CreateStructure ¶
func (s *SchemaOps) CreateStructure(ctx context.Context, model *unifiedmodel.UnifiedModel) error
CreateStructure creates Databricks structure from a UnifiedModel.
func (*SchemaOps) DiscoverSchema ¶
func (s *SchemaOps) DiscoverSchema(ctx context.Context) (*unifiedmodel.UnifiedModel, error)
DiscoverSchema retrieves the schema of Databricks database (Delta tables).
func (*SchemaOps) GetTableSchema ¶
func (s *SchemaOps) GetTableSchema(ctx context.Context, tableName string) (*unifiedmodel.Table, error)
GetTableSchema retrieves the schema for a specific table.