package main import ( "context" "encoding/json" "flag" "fmt" "log" "os" "os/signal" "strings" "syscall" "time" "github.com/IBM/sarama" "github.com/neo4j/neo4j-go-driver/v5/neo4j" ) type Config struct { KafkaBrokers string KafkaTopic string KafkaGroup string Neo4jURI string Neo4jUser string Neo4jPass string } func main() { var cfg Config flag.StringVar(&cfg.KafkaBrokers, "kafka-brokers", envOr("KAFKA_BROKERS", "localhost:9092"), "Kafka brokers") flag.StringVar(&cfg.KafkaTopic, "kafka-topic", envOr("KAFKA_TOPIC", "muyu_crm.public.crm_tenant_relationship"), "CDC topic") flag.StringVar(&cfg.KafkaGroup, "kafka-group", envOr("KAFKA_GROUP", "graphsync"), "consumer group") flag.StringVar(&cfg.Neo4jURI, "neo4j-uri", envOr("NEO4J_URI", "bolt://localhost:7687"), "Neo4j bolt URI") flag.StringVar(&cfg.Neo4jUser, "neo4j-user", envOr("NEO4J_USER", "neo4j"), "Neo4j user") flag.StringVar(&cfg.Neo4jPass, "neo4j-pass", envOr("NEO4J_PASS", ""), "Neo4j password") flag.Parse() ctx, cancel := context.WithCancel(context.Background()) defer cancel() driver, err := neo4j.NewDriverWithContext(cfg.Neo4jURI, neo4j.BasicAuth(cfg.Neo4jUser, cfg.Neo4jPass, "")) if err != nil { log.Fatalf("neo4j connect failed: %v", err) } defer driver.Close(ctx) saramaCfg := sarama.NewConfig() saramaCfg.Consumer.Group.Rebalance.GroupStrategies = []sarama.BalanceStrategy{sarama.NewBalanceStrategyRoundRobin()} saramaCfg.Consumer.Offsets.Initial = sarama.OffsetOldest group, err := sarama.NewConsumerGroup(strings.Split(cfg.KafkaBrokers, ","), cfg.KafkaGroup, saramaCfg) if err != nil { log.Fatalf("kafka consumer group failed: %v", err) } defer group.Close() handler := &syncHandler{driver: driver} go func() { for { if err := group.Consume(ctx, []string{cfg.KafkaTopic}, handler); err != nil { log.Printf("consume error: %v, retrying in 5s...", err) time.Sleep(5 * time.Second) } if ctx.Err() != nil { return } } }() log.Printf("GraphSyncWorker started: brokers=%s topic=%s neo4j=%s", cfg.KafkaBrokers, cfg.KafkaTopic, cfg.Neo4jURI) sig := make(chan os.Signal, 1) signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM) <-sig log.Println("shutting down...") cancel() } type syncHandler struct { driver neo4j.DriverWithContext } func (h *syncHandler) Setup(_ sarama.ConsumerGroupSession) error { return nil } func (h *syncHandler) Cleanup(_ sarama.ConsumerGroupSession) error { return nil } func (h *syncHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { for msg := range claim.Messages() { if err := h.processMessage(session.Context(), msg.Value); err != nil { log.Printf("process msg offset=%d err: %v", msg.Offset, err) } session.MarkMessage(msg, "") } return nil } type cdcEvent struct { Op string `json:"op"` After map[string]interface{} `json:"after"` } func (h *syncHandler) processMessage(ctx context.Context, data []byte) error { var evt cdcEvent if err := json.Unmarshal(data, &evt); err != nil { return fmt.Errorf("unmarshal: %w", err) } if evt.After == nil { return nil } session := h.driver.NewSession(ctx, neo4j.SessionConfig{AccessMode: neo4j.AccessModeWrite}) defer session.Close(ctx) fromTenant, _ := evt.After["from_tenant_id"].(string) toTenant, _ := evt.After["to_tenant_id"].(string) relType, _ := evt.After["relation_type"].(string) relID, _ := evt.After["relationship_id"].(string) if fromTenant == "" || toTenant == "" { return nil } cypher := ` MERGE (a:Tenant {tenantId: $from}) MERGE (b:Tenant {tenantId: $to}) MERGE (a)-[r:SUPPLIES {relationshipId: $relId}]->(b) SET r.relationType = $relType, r.updatedAt = datetime() ` _, err := session.Run(ctx, cypher, map[string]interface{}{ "from": fromTenant, "to": toTenant, "relId": relID, "relType": relType, }) if err != nil { return fmt.Errorf("neo4j exec: %w", err) } log.Printf("synced: %s -> %s [%s] id=%s", fromTenant, toTenant, relType, relID) return nil } func envOr(key, fallback string) string { if v := os.Getenv(key); v != "" { return v } return fallback }