150 lines
4.1 KiB
Go
150 lines
4.1 KiB
Go
|
|
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", "muyu2026"), "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
|
||
|
|
}
|