This is an automated email from the ASF dual-hosted git repository.
alfredlu pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 756f9cb [INLONG-2066] Confirm the message correctly in Go SDK (#2067)
756f9cb is described below
commit 756f9cb591056383b3fdd4c2ce85b936fffc6466
Author: Zijie Lu <[email protected]>
AuthorDate: Tue Dec 28 15:22:20 2021 +0800
[INLONG-2066] Confirm the message correctly in Go SDK (#2067)
Signed-off-by: Zijie Lu <[email protected]>
---
.../tubemq-client-go/client/consumer_impl.go | 13 +++---
.../tubemq-client-go/config/config.go | 34 +++++++++++----
.../tubemq-client-go/config/config_test.go | 49 ++++++++++++++++++++++
.../tubemq-client-go/tdmsg/td_msg_test.go | 2 +-
4 files changed, 84 insertions(+), 14 deletions(-)
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/consumer_impl.go
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/consumer_impl.go
index bc26097..ca9cca0 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/consumer_impl.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/client/consumer_impl.go
@@ -280,7 +280,8 @@ func (c *consumer) Confirm(confirmContext string, consumed
bool) (*ConsumerResul
return nil, errs.New(errs.RetErrConfirmTimeout, "Not found the
partition by confirm_context!")
}
- rsp, err := c.sendConfirmReq2Broker(partition)
+ defer c.rmtDataCache.ReleasePartition(true,
c.subInfo.IsFiltered(topic), confirmContext, consumed)
+ rsp, err := c.sendConfirmReq2Broker(partition, consumed)
if err != nil {
log.Infof("[CONSUMER]Confirm error %s", err.Error())
return nil, err
@@ -290,7 +291,8 @@ func (c *consumer) Confirm(confirmContext string, consumed
bool) (*ConsumerResul
BrokerHost: partition.GetBroker().GetHost(),
PartitionID: uint32(partition.GetPartitionID()),
PartitionKey: partition.GetPartitionKey(),
- CurrOffset: util.InvalidValue,
+ CurrOffset: rsp.GetCurrOffset(),
+ MaxOffset: rsp.GetMaxOffset(),
}
cs := &ConsumerResult{
TopicName: partition.GetTopic(),
@@ -299,13 +301,11 @@ func (c *consumer) Confirm(confirmContext string,
consumed bool) (*ConsumerResul
if !rsp.GetSuccess() {
return cs, errs.New(rsp.GetErrCode(), rsp.GetErrMsg())
}
- currOffset := rsp.GetCurrOffset()
- c.rmtDataCache.BookPartitionInfo(partitionKey, currOffset,
util.InvalidValue)
- err = c.rmtDataCache.ReleasePartition(true,
c.subInfo.IsFiltered(topic), confirmContext, consumed)
+ c.rmtDataCache.BookPartitionInfo(partitionKey, rsp.GetCurrOffset(),
rsp.GetMaxOffset())
return cs, err
}
-func (c *consumer) sendConfirmReq2Broker(partition *metadata.Partition)
(*protocol.CommitOffsetResponseB2C, error) {
+func (c *consumer) sendConfirmReq2Broker(partition *metadata.Partition,
consumed bool) (*protocol.CommitOffsetResponseB2C, error) {
m := &metadata.Metadata{}
node := &metadata.Node{}
node.SetHost(util.GetLocalHost())
@@ -313,6 +313,7 @@ func (c *consumer) sendConfirmReq2Broker(partition
*metadata.Partition) (*protoc
m.SetNode(node)
sub := &metadata.SubscribeInfo{}
sub.SetGroup(c.config.Consumer.Group)
+ partition.SetLastConsumed(consumed)
sub.SetPartition(partition)
m.SetSubscribeInfo(sub)
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config.go
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config.go
index bd54505..f419d1a 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config.go
@@ -32,13 +32,16 @@ import (
)
const (
- MaxRPCTimeout = 300000 * time.Millisecond
- MinRPCTimeout = 8000 * time.Millisecond
- MaxSessionKeyLen = 1024
- MaxGroupLen = 1024
- MaxTopicLen = 64
- MaxFilterLen = 256
- MaxFilterItemCount = 500
+ MaxRPCTimeout = 300000 * time.Millisecond
+ MinRPCTimeout = 8000 * time.Millisecond
+ MaxSessionKeyLen = 1024
+ MaxGroupLen = 1024
+ MaxTopicLen = 64
+ MaxFilterLen = 256
+ MaxFilterItemCount = 500
+ ConsumeFromFirstOffset = -1
+ ConsumeFromLatestOffset = 0
+ ConsumeFromMaxOffsetAlways = 1
)
// Config defines multiple configuration options.
@@ -196,6 +199,10 @@ func (c *Config) ValidateConsumer() error {
}
}
+ if err := c.validateConsumePosition(); err != nil {
+ return err
+ }
+
err := c.validateMaster()
if err != nil {
return errs.New(errs.RetInvalidConfig, err.Error())
@@ -302,6 +309,13 @@ func (c *Config) validateGroup(group string) error {
return nil
}
+func (c *Config) validateConsumePosition() error {
+ if c.Consumer.ConsumePosition < ConsumeFromFirstOffset ||
c.Consumer.ConsumePosition > ConsumeFromMaxOffsetAlways {
+ return errs.New(errs.RetInvalidConfig,
fmt.Sprintf("consumePosition should be only in (-1, 0, 1), while %d is passed",
c.Consumer.ConsumePosition))
+ }
+ return nil
+}
+
// ParseAddress parses the address to user-defined config.
func ParseAddress(address string) (config *Config, err error) {
c := NewDefaultConfig()
@@ -545,3 +559,9 @@ func WithBoundConsume(sessionKey string, sourceCount int,
selectBig bool, partOf
c.Consumer.PartitionOffset = partOffset
}
}
+
+func WithConsumePosition(consumePosition int) Option {
+ return func(c *Config) {
+ c.Consumer.ConsumePosition = consumePosition
+ }
+}
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config_test.go
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config_test.go
index c284935..9d5e910 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config_test.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/config/config_test.go
@@ -67,6 +67,40 @@ func TestParseAddress(t *testing.T) {
address = "127.0.0.1:9092,127.0.0.1:9093?filters=12312323&filters=1212"
_, err = ParseAddress(address)
assert.NotNil(t, err)
+
+ address =
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=-1"
+ c, err = ParseAddress(address)
+ assert.Nil(t, err)
+ err = c.ValidateConsumer()
+ assert.Nil(t, err)
+
+ address =
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=1"
+ c, err = ParseAddress(address)
+ assert.Nil(t, err)
+ err = c.ValidateConsumer()
+ assert.Nil(t, err)
+
+ address =
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=0"
+ c, err = ParseAddress(address)
+ assert.Nil(t, err)
+ err = c.ValidateConsumer()
+ assert.Nil(t, err)
+
+ address =
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=-2"
+ c, err = ParseAddress(address)
+ assert.Nil(t, err)
+ err = c.ValidateConsumer()
+ assert.NotNil(t, err)
+
+ address =
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=2"
+ c, err = ParseAddress(address)
+ assert.Nil(t, err)
+ err = c.ValidateConsumer()
+ assert.NotNil(t, err)
+
+ address =
"127.0.0.1:9092,127.0.0.1:9093?topic=Topic&group=group&filters=12312323&consumePosition=a"
+ c, err = ParseAddress(address)
+ assert.NotNil(t, err)
}
func TestValidateGroup(t *testing.T) {
@@ -131,4 +165,19 @@ func TestValidateConsumer(t *testing.T) {
partitionOffset = map[string]int64{"181895251:topic1:1": 0,
"181895251:topic2:2": 10}
WithBoundConsume("11", 0, true, partitionOffset)(c)
assert.Nil(t, c.ValidateConsumer())
+
+ WithConsumePosition(2)(c)
+ assert.NotNil(t, c.ValidateConsumer())
+
+ WithConsumePosition(-2)(c)
+ assert.NotNil(t, c.ValidateConsumer())
+
+ WithConsumePosition(-1)(c)
+ assert.Nil(t, c.ValidateConsumer())
+
+ WithConsumePosition(0)(c)
+ assert.Nil(t, c.ValidateConsumer())
+
+ WithConsumePosition(1)(c)
+ assert.Nil(t, c.ValidateConsumer())
}
diff --git
a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/tdmsg/td_msg_test.go
b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/tdmsg/td_msg_test.go
index 1e8408d..70005f6 100644
--- a/inlong-tubemq/tubemq-client-twins/tubemq-client-go/tdmsg/td_msg_test.go
+++ b/inlong-tubemq/tubemq-client-twins/tubemq-client-go/tdmsg/td_msg_test.go
@@ -34,7 +34,7 @@ func TestTDMsgV4(t *testing.T) {
}
func TestTDMsgV1(t *testing.T) {
- b :=
[]byte{15,1,0,0,1,125,71,98,161,138,0,0,0,1,0,206,100,116,61,49,54,51,55,53,56,48,49,56,53,57,56,55,38,109,115,103,85,85,73,68,61,100,98,52,97,51,101,51,100,45,50,100,101,55,45,52,99,102,102,45,56,54,97,101,45,98,53,55,52,55,101,57,49,98,51,101,53,38,99,110,116,61,49,38,109,116,61,112,98,38,78,111,100,101,73,80,61,49,49,46,49,53,52,46,50,48,57,46,49,55,57,38,114,116,61,49,54,51,55,53,56,48,49,56,53,57,57,52,38,109,61,57,38,116,105,100,61,116,95,115,110,103,95,103,100,116,95,117,110
[...]
+ b := []byte{15, 1, 0, 0, 1, 125, 71, 98, 161, 138, 0, 0, 0, 1, 0, 206,
100, 116, 61, 49, 54, 51, 55, 53, 56, 48, 49, 56, 53, 57, 56, 55, 38, 109, 115,
103, 85, 85, 73, 68, 61, 100, 98, 52, 97, 51, 101, 51, 100, 45, 50, 100, 101,
55, 45, 52, 99, 102, 102, 45, 56, 54, 97, 101, 45, 98, 53, 55, 52, 55, 101, 57,
49, 98, 51, 101, 53, 38, 99, 110, 116, 61, 49, 38, 109, 116, 61, 112, 98, 38,
78, 111, 100, 101, 73, 80, 61, 49, 49, 46, 49, 53, 52, 46, 50, 48, 57, 46, 49,
55, 57, 38, 114, 116, 61, [...]
tm, err := New(b)
assert.Equal(t, uint64(1637580185994), tm.CreateTime)
assert.Equal(t, int32(1), tm.Version)