Skip to content

Commit

Permalink
Add schema support to Reader
Browse files Browse the repository at this point in the history
  • Loading branch information
ZiyaoWei committed Mar 8, 2022
1 parent 1be07f9 commit 81e315c
Show file tree
Hide file tree
Showing 3 changed files with 47 additions and 0 deletions.
3 changes: 3 additions & 0 deletions pulsar/reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,9 @@ type ReaderOptions struct {

// Decryption represents the encryption related fields required by the reader to decrypt a message.
Decryption *MessageDecryptionInfo

// Schema represents the schema implementation.
Schema Schema
}

// Reader can be used to scan through all the messages currently available in a topic.
Expand Down
1 change: 1 addition & 0 deletions pulsar/reader_impl.go
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,7 @@ func newReader(client *client, options ReaderOptions) (Reader, error) {
nackRedeliveryDelay: defaultNackRedeliveryDelay,
replicateSubscriptionState: false,
decryption: options.Decryption,
schema: options.Schema,
}

reader := &reader{
Expand Down
43 changes: 43 additions & 0 deletions pulsar/reader_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -710,3 +710,46 @@ func TestProducerReaderRSAEncryption(t *testing.T) {
assert.Equal(t, []byte(expectMsg), msg.Payload())
}
}

func TestReaderWithSchema(t *testing.T) {
client, err := NewClient(ClientOptions{
URL: lookupURL,
})

assert.Nil(t, err)
defer client.Close()

topic := newTopicName()
schema := NewStringSchema(nil)

// create producer
producer, err := client.CreateProducer(ProducerOptions{
Topic: topic,
Schema: schema,
})
assert.Nil(t, err)
defer producer.Close()

value := "hello pulsar"
_, err = producer.Send(context.Background(), &ProducerMessage{
Value: value,
})
assert.Nil(t, err)

// create reader
reader, err := client.CreateReader(ReaderOptions{
Topic: topic,
StartMessageID: EarliestMessageID(),
Schema: schema,
})
assert.Nil(t, err)
defer reader.Close()

msg, err := reader.Next(context.Background())
assert.NoError(t, err)

var res *string
err = msg.GetSchemaValue(&res)
assert.Nil(t, err)
assert.Equal(t, *res, value)
}

0 comments on commit 81e315c

Please sign in to comment.