This repository has been archived by the owner on Sep 15, 2023. It is now read-only.
-
-
Notifications
You must be signed in to change notification settings - Fork 36
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge pull request #144 from askuy/feature/kafkatrace
kafka custom trace
- Loading branch information
Showing
9 changed files
with
173 additions
and
53 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,5 @@ | ||
run:export EGO_DEBUG=true | ||
run:export EGO_LOG_EXTRA_KEYS=X-Ego-Uid | ||
run: | ||
go run main.go --config=config.toml | ||
|
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,13 @@ | ||
[kafka] | ||
debug = true | ||
EnableAccessInterceptor = true | ||
EnableAccessInterceptorReq = true | ||
EnableAccessInterceptorRes = true | ||
brokers = ["10.8.0.1:9092"] | ||
[kafka.client] | ||
timeout = "3s" | ||
[kafka.producers.p1] # 定义了名字为p1的producer | ||
topic = "test" # 指定生产消息的topic | ||
[kafka.consumers.c1] # 定义了名字为c1的consumer | ||
topic = "test" # 指定消费的topic | ||
groupID = "group-1" # 如果配置了groupID,将初始化为consumerGroup |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,62 @@ | ||
package main | ||
|
||
import ( | ||
"context" | ||
"fmt" | ||
"log" | ||
|
||
"github.com/gotomicro/ego" | ||
"github.com/gotomicro/ego-component/ekafka" | ||
"github.com/gotomicro/ego/core/transport" | ||
) | ||
|
||
// export EGO_DEBUG=true | ||
func main() { | ||
ego.New().Invoker(func() error { | ||
ctx := context.Background() | ||
ctx = transport.WithValue(ctx, "X-Ego-Uid", 9527) | ||
// 初始化ekafka组件 | ||
cmp := ekafka.Load("kafka").Build() | ||
// 使用p1生产者生产消息 | ||
produce(ctx, cmp.Producer("p1")) | ||
// 使用c1消费者消费消息 | ||
consume(cmp.Consumer("c1")) | ||
return nil | ||
}).Run() | ||
|
||
} | ||
|
||
// produce 生产消息 | ||
func produce(ctx context.Context, w *ekafka.Producer) { | ||
// 生产3条消息 | ||
ctx = context.WithValue(ctx, "hello", "world") | ||
err := w.WriteMessages(ctx, | ||
&ekafka.Message{Key: []byte("Key-A"), Value: []byte("Hellohahah World!22222")}, | ||
) | ||
if err != nil { | ||
log.Fatal("failed to write messages:", err) | ||
} | ||
if err := w.Close(); err != nil { | ||
log.Fatal("failed to close writer:", err) | ||
} | ||
} | ||
|
||
// consume 使用consumer/consumerGroup消费消息 | ||
func consume(r *ekafka.Consumer) { | ||
ctx := context.Background() | ||
for { | ||
// ReadMessage 再收到下一个Message时,会阻塞 | ||
msg, ctxOutput, err := r.ReadMessage(ctx) | ||
if err != nil { | ||
panic("could not read message " + err.Error()) | ||
} | ||
|
||
// 打印消息 | ||
fmt.Println("received headers: ", msg.Headers) | ||
fmt.Println("received: ", string(msg.Value)) | ||
err = r.CommitMessages(ctxOutput, &msg) | ||
if err != nil { | ||
log.Printf("fail to commit msg:%v", err) | ||
} | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters