databricks

package
v0.0.0-...-ba04163 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Nov 25, 2025 License: AGPL-3.0 Imports: 9 Imported by: 0

Documentation

Index

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

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

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) Close

func (c *Connection) Close() error

Close closes the connection.

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) ID

func (c *Connection) ID() string

ID returns the connection identifier.

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

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

func (d *DataOps) ExecuteCountQuery(ctx context.Context, query string) (int64, error)

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

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.

func (*DataOps) Upsert

func (d *DataOps) Upsert(ctx context.Context, table string, data []map[string]interface{}, uniqueColumns []string) (int64, error)

Upsert performs upsert operation in Databricks using MERGE.

func (*DataOps) Wipe

func (d *DataOps) Wipe(ctx context.Context) error

Wipe deletes all rows from all tables in the database.

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.

func (*DatabricksClient) Ping

func (c *DatabricksClient) Ping(ctx context.Context) error

Ping tests the Databricks connection.

type InstanceConnection

type InstanceConnection struct {
	// contains filtered or unexported fields
}

InstanceConnection implements adapter.InstanceConnection for Databricks.

func (*InstanceConnection) Adapter

Adapter returns the database adapter.

func (*InstanceConnection) Close

func (ic *InstanceConnection) Close() error

Close closes the connection.

func (*InstanceConnection) Config

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

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

func (m *MetadataOps) ExecuteCommand(ctx context.Context, command string) ([]byte, error)

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

func (r *ReplicationOps) ApplyCDCEvent(ctx context.Context, event *adapter.CDCEvent) error

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

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.

func (*SchemaOps) ListTables

func (s *SchemaOps) ListTables(ctx context.Context) ([]string, error)

ListTables lists all tables in the database.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL