go-scylla-cdc
A Go wrapper for ScyllaDB Change Data Capture (CDC) that simplifies consumer implementation with automatic progress reporting and graceful shutdown handling.
Installation
go get github.com/gokpm/go-scylla-cdc
Usage
package main
import (
"context"
"log"
"time"
"github.com/gocql/gocql"
"github.com/gokpm/go-scylla-cdc"
scyllacdc "github.com/scylladb/scylla-cdc-go"
)
func main() {
// Create ScyllaDB cluster connection
cluster := gocql.NewCluster("127.0.0.1")
cluster.Keyspace = "your_keyspace"
cluster.Consistency = gocql.LocalQuorum
cluster.Timeout = 10 * time.Second
cluster.ConnectTimeout = 10 * time.Second
// Create session
session, err := cluster.CreateSession()
if err != nil {
log.Fatal("Failed to create session:", err)
}
defer session.Close()
// Create CDC reader with advanced configuration
readerConfig := &scyllacdc.ReaderConfig{
Session: session,
TableNames: []string{"your_table"},
ChangeConsumerFactory: nil, // Will be set below
Advanced: scyllacdc.AdvancedReaderConfig{
ChangeAgeLimit: 30 * time.Minute,
ConfidenceWindowSize: 10 * time.Second,
PostNonEmptyQueryDelay: 1 * time.Second,
PostEmptyQueryDelay: 5 * time.Second,
PostFailedQueryDelay: 30 * time.Second,
QueryTimeWindowSize: 60 * time.Second,
},
}
// Create progress manager
progressManager, err := scyllacdc.NewTableBackedProgressManager(
session,
"cdc_progress", // progress table name
"my_app", // application name
)
if err != nil {
log.Fatal("Failed to create progress manager:", err)
}
// Configure your CDC consumer
config := cdc.Config{
CommitInterval: 5 * time.Second,
ShutdownTimeout: 30 * time.Second,
Consumer: func(ctx context.Context, change scyllacdc.Change) error {
// Process your CDC change here
log.Printf("Received change: %+v", change)
return nil
},
}
// Create factory and set it in reader config
factory := cdc.New(config)
readerConfig.ChangeConsumerFactory = factory
// Create and start CDC reader
reader, err := scyllacdc.NewReader(context.Background(), *readerConfig)
if err != nil {
log.Fatal("Failed to create CDC reader:", err)
}
// Start reading changes
err = reader.Run(context.Background())
if err != nil {
log.Fatal("CDC reader failed:", err)
}
}
Features
- Automatic periodic progress reporting with table-backed persistence
- Configurable commit intervals and advanced reader settings
- Graceful shutdown with timeout
- Error handling and logging integration
- Full ScyllaDB cluster connection and session management
- Compatible with
github.com/scylladb/scylla-cdc-go
Configuration
CDC Config
CommitInterval: How often to commit progress
Consumer: Your change processing function
ShutdownTimeout: Maximum time to wait during shutdown
Advanced Reader Config
ChangeAgeLimit: Maximum age of changes to process
ConfidenceWindowSize: Time window for confidence calculations
- Query delay settings for different scenarios
- Generation and time window configurations
Dependencies
go get github.com/gocql/gocql
go get github.com/scylladb/scylla-cdc-go