micro--go-micro
7e1bba2baf
* Re-add events package * run redis as a dep * remove redis events * fix: data race on event subscriber * fix: data race in tests * fix: store errors * fix: lint issues * feat: default stream * Update file.go --------- Co-authored-by: Brian Ketelsen <bketelsen@gmail.com>
49 行
816 B
Markdown
49 行
816 B
Markdown
# NATS JetStream
|
|
|
|
This plugin uses NATS with JetStream to send and receive events.
|
|
|
|
## Create a stream
|
|
|
|
```go
|
|
ev, err := natsjs.NewStream(
|
|
natsjs.Address("nats://10.0.1.46:4222"),
|
|
natsjs.MaxAge(24*160*time.Minute),
|
|
)
|
|
```
|
|
|
|
## Consume a stream
|
|
|
|
```go
|
|
ee, err := events.Consume("test",
|
|
events.WithAutoAck(false, time.Second*30),
|
|
events.WithGroup("testgroup"),
|
|
)
|
|
if err != nil {
|
|
panic(err)
|
|
}
|
|
go func() {
|
|
for {
|
|
msg := <-ee
|
|
// Process the message
|
|
logger.Info("Received message:", string(msg.Payload))
|
|
err := msg.Ack()
|
|
if err != nil {
|
|
logger.Error("Error acknowledging message:", err)
|
|
} else {
|
|
logger.Info("Message acknowledged")
|
|
}
|
|
}
|
|
}()
|
|
|
|
```
|
|
|
|
## Publish an Event to the stream
|
|
|
|
```go
|
|
err = ev.Publish("test", []byte("hello world"))
|
|
if err != nil {
|
|
panic(err)
|
|
}
|
|
```
|
|
|