cdc

package module
v0.0.0-...-e298809 Latest Latest
Warning

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

Go to latest
Published: Jul 1, 2025 License: MIT Imports: 4 Imported by: 0

README

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

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func New

func New(config *Config) *factory

Types

type Config

type Config struct {
	CommitInterval  time.Duration
	Consumer        func(context.Context, cdc.Change) error
	ShutdownTimeout time.Duration
}

Jump to

Keyboard shortcuts

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