INNER CODE UNIT · Go

pushSubscribeOnOrder

shijuvar/go-distributed-sys · eventstream/querymodelworker/main.go:36

func pushSubscribeOnOrder(js nats.JetStreamContext) {
	// Create durable push consumer
	js.QueueSubscribe(subscribeSubject, queueGroup, func(msg *nats.Msg) {
		msg.Ack()
		var order ordermodel.Order
		// Unmarshal JSON that represents the Order data
		err := json.Unmarshal(msg.Data, &order)
		if err != nil {
			log.Print(err)
			return
		}
		log.Printf("Message subscribed on subject:%s, from:%s,  data:%v", subscribeSubject, clientID, order)
		orderDB, _ := sqldb.NewOrdersDB()
		repository, _ := ordersyncrepository.New(orderDB.DB)
		// Sync query model with event data
		if err := repository.CreateOrder(context.Background(), order); err != nil {
			log.Printf("Error while replicating the query model %+v", err)
		}

View source record →

📰 Research Paper
Loading…
⏳ Fetching content…