Skip to content

Commit 2ba8285

Browse files
authored
Merge pull request #600 from chiemezie1/feat/outbox-kafka-publisher
feat: add Kafka publisher for outbox dispatcher
2 parents 3490efb + 5e6ca1c commit 2ba8285

7 files changed

Lines changed: 480 additions & 1 deletion

File tree

docs/outbox-pattern.md

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -324,6 +324,17 @@ The system automatically creates the outbox table on startup. For production dep
324324
3. **Compression**: Event payload compression for large events
325325
4. **Batch Publishing**: Batch multiple events to external systems
326326

327+
### Kafka Publisher Option
328+
329+
Outbox dispatch can now route selected topics to Kafka instead of HTTP. Configure it with the following environment variables:
330+
331+
- `OUTBOX_KAFKA_BROKERS`: comma-separated broker list
332+
- `OUTBOX_KAFKA_TOPIC_MAP`: comma-separated `event_type=topic` mappings
333+
- `OUTBOX_KAFKA_ACKS`: default Kafka ack level (`0`, `1`, or `all`)
334+
- `OUTBOX_KAFKA_TOPIC_ACKS`: per-topic overrides such as `billing-events=all`
335+
336+
When the service is configured with `PublisherType: "kafka"`, the dispatcher uses the Kafka publisher for mapped topics and preserves the existing HTTP path for unconfigured topics. Kafka publish latency and errors are emitted through the `outbox_kafka_produce_latency_seconds` and `outbox_kafka_errors_total` metrics.
337+
327338
## Examples
328339

329340
### Example Domain Event

go.mod

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -123,6 +123,7 @@ require (
123123
github.com/opencontainers/image-spec v1.1.1 // indirect
124124
github.com/pelletier/go-toml/v2 v2.2.4 // indirect
125125
github.com/perimeterx/marshmallow v1.1.5 // indirect
126+
github.com/pierrec/lz4/v4 v4.1.15 // indirect
126127
github.com/pkg/errors v0.9.1 // indirect
127128
github.com/pmezard/go-difflib v1.0.0 // indirect
128129
github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect
@@ -134,6 +135,7 @@ require (
134135
github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9 // indirect
135136
github.com/redis/go-redis/v9 v9.7.1 // indirect
136137
github.com/segmentio/asm v1.2.1 // indirect
138+
github.com/segmentio/kafka-go v0.4.47 // indirect
137139
github.com/shirou/gopsutil/v4 v4.26.2 // indirect
138140
github.com/stretchr/objx v0.5.2 // indirect
139141
github.com/tchap/go-patricia/v2 v2.3.3 // indirect

go.sum

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -261,6 +261,8 @@ github.com/pelletier/go-toml/v2 v2.2.4 h1:mye9XuhQ6gvn5h28+VilKrrPoQVanw5PMw/TB0
261261
github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY=
262262
github.com/perimeterx/marshmallow v1.1.5 h1:a2LALqQ1BlHM8PZblsDdidgv1mWi1DgC2UmX50IvK2s=
263263
github.com/perimeterx/marshmallow v1.1.5/go.mod h1:dsXbUu8CRzfYP5a87xpp0xq9S3u0Vchtcl8we9tYaXw=
264+
github.com/pierrec/lz4/v4 v4.1.15 h1:MO0/ucJhngq7299dKLwIMtgTfbkoSPF6AoMYDd8Q4q0=
265+
github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4=
264266
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
265267
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
266268
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
@@ -298,6 +300,8 @@ github.com/segmentio/asm v1.2.0 h1:9BQrFxC+YOHJlTlHGkTrFWf59nbL3XnCoFLTwDCI7ys=
298300
github.com/segmentio/asm v1.2.0/go.mod h1:BqMnlJP91P8d+4ibuonYZw9mfnzI9HfxselHZr5aAcs=
299301
github.com/segmentio/asm v1.2.1 h1:DTNbBqs57ioxAD4PrArqftgypG4/qNpXoJx8TVXxPR0=
300302
github.com/segmentio/asm v1.2.1/go.mod h1:BqMnlJP91P8d+4ibuonYZw9mfnzI9HfxselHZr5aAcs=
303+
github.com/segmentio/kafka-go v0.4.47 h1:IqziR4pA3vrZq7YdRxaT3w1/5fvIH5qpCwstUanQQB0=
304+
github.com/segmentio/kafka-go v0.4.47/go.mod h1:HjF6XbOKh0Pjlkr5GVZxt6CsjjwnmhVOfURM5KMd8qg=
301305
github.com/shirou/gopsutil/v4 v4.26.2 h1:X8i6sicvUFih4BmYIGT1m2wwgw2VG9YgrDTi7cIRGUI=
302306
github.com/shirou/gopsutil/v4 v4.26.2/go.mod h1:LZ6ewCSkBqUpvSOf+LsTGnRinC6iaNUNMGBtDkJBaLQ=
303307
github.com/shopspring/decimal v1.4.0 h1:bxl37RwXBklmTi0C79JfXCEBD1cqqHt0bbgBAGFp81k=
@@ -344,6 +348,9 @@ github.com/vektah/gqlparser/v2 v2.5.34 h1:MEea5P0qhdcqfBL45ghKE+qr9laidVHTMHjav5
344348
github.com/vektah/gqlparser/v2 v2.5.34/go.mod h1:mFdHLGCio7OGX1fby9ZjTW6FN+qxgmbnBcRIeeScE5s=
345349
github.com/woodsbury/decimal128 v1.3.0 h1:8pffMNWIlC0O5vbyHWFZAt5yWvWcrHA+3ovIIjVWss0=
346350
github.com/woodsbury/decimal128 v1.3.0/go.mod h1:C5UTmyTjW3JftjUFzOVhC20BEQa2a4ZKOB5I6Zjb+ds=
351+
github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI=
352+
github.com/xdg-go/scram v1.1.2/go.mod h1:RT/sEzTbU5y00aCK8UOx6R7YryM0iF1N2MOmC3kKLN4=
353+
github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM=
347354
github.com/xeipuuv/gojsonpointer v0.0.0-20180127040702-4e3ac2762d5f/go.mod h1:N2zxlSyiKSe5eX1tZViRH5QA0qijqEDrYZiPEAiq3wU=
348355
github.com/xeipuuv/gojsonpointer v0.0.0-20190905194746-02993c407bfb h1:zGWFAtiMcyryUHoUjUJX0/lt1H2+i2Ka2n+D3DImSNo=
349356
github.com/xeipuuv/gojsonpointer v0.0.0-20190905194746-02993c407bfb/go.mod h1:N2zxlSyiKSe5eX1tZViRH5QA0qijqEDrYZiPEAiq3wU=
@@ -353,6 +360,7 @@ github.com/xeipuuv/gojsonschema v1.2.0 h1:LhYJRs+L4fBtjZUfuSZIKGeVu0QRy8e5Xi7D17
353360
github.com/xeipuuv/gojsonschema v1.2.0/go.mod h1:anYRn/JVcOK2ZgGU+IjEV4nwlhoK5sQluxsYJ78Id3Y=
354361
github.com/yashtewari/glob-intersection v0.2.0 h1:8iuHdN88yYuCzCdjt0gDe+6bAhUwBeEWqThExu54RFg=
355362
github.com/yashtewari/glob-intersection v0.2.0/go.mod h1:LK7pIC3piUjovexikBbJ26Yml7g8xa5bsjfx2v1fwok=
363+
github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY=
356364
github.com/yuin/gopher-lua v1.1.1 h1:kYKnWBjvbNP4XLT3+bPEwAXJx262OhaHDWDVOPjL46M=
357365
github.com/yuin/gopher-lua v1.1.1/go.mod h1:GBR0iDaNXjAgGg9zfCvksxSRnQx76gclCIb7kdAd1Pw=
358366
github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0=
@@ -432,12 +440,26 @@ go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=
432440
go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
433441
golang.org/x/arch v0.24.0 h1:qlJ3M9upxvFfwRM51tTg3Yl+8CP9vCC1E7vlFpgv99Y=
434442
golang.org/x/arch v0.24.0/go.mod h1:dNHoOeKiyja7GTvF9NJS1l3Z2yntpQNzgrjh1cU103A=
443+
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
444+
golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc=
445+
golang.org/x/crypto v0.14.0/go.mod h1:MVFd36DqK4CsrnJYDkBA3VC4m2GkXAM0PvzMCn4JQf4=
435446
golang.org/x/crypto v0.52.0 h1:RMs7fP2rXdep0CftQlK8Uf+kibLm7qkCcradZWYz988=
436447
golang.org/x/crypto v0.52.0/go.mod h1:1QgfPxDqh0T2M/elOJtp9RvuR95kVjir0e6/BvEmGbc=
448+
golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4=
449+
golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
437450
golang.org/x/mod v0.36.0 h1:JJjpVx6myfUsUdAzZuOSTTmRE0PfZeNWzzvKrP7amb4=
438451
golang.org/x/mod v0.36.0/go.mod h1:moc6ELqsWcOw5Ef3xVprK5ul/MvtVvkIXLziUOICjUQ=
452+
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
453+
golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg=
454+
golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c=
455+
golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs=
456+
golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg=
457+
golang.org/x/net v0.17.0/go.mod h1:NxSsAGuq816PNPmqtQdLE42eU2Fs7NoRIZrHJAlaCOE=
439458
golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8=
440459
golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww=
460+
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
461+
golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
462+
golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
441463
golang.org/x/crypto v0.48.0 h1:/VRzVqiRSggnhY7gNRxPauEQ5Drw9haKdM0jqfcCFts=
442464
golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos=
443465
golang.org/x/crypto v0.52.0 h1:RMs7fP2rXdep0CftQlK8Uf+kibLm7qkCcradZWYz988=
@@ -450,20 +472,45 @@ golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4=
450472
golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI=
451473
golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM=
452474
golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
475+
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
453476
golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
477+
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
454478
golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
479+
golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
455480
golang.org/x/sys v0.0.0-20210616094352-59db8d763f22/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
481+
golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
482+
golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
483+
golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
456484
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
485+
golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
486+
golang.org/x/sys v0.13.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
457487
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY=
458488
golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
489+
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
490+
golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8=
491+
golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k=
492+
golang.org/x/term v0.8.0/go.mod h1:xPskH00ivmX89bAKVGSKKtLOWNx2+17Eiy94tnKShWo=
493+
golang.org/x/term v0.13.0/go.mod h1:LTmsnFJwVN6bCy1rVCoS+qHT1HhALEFxKncY3WNNh4U=
459494
golang.org/x/term v0.43.0 h1:S4RLU2sB31O/NCl+zFN9Aru9A/Cq2aqKpTZJ6B+DwT4=
460495
golang.org/x/term v0.43.0/go.mod h1:lrhlHNdQJHO+1qVYiHfFKVuVioJIheAc3fBSMFYEIsk=
496+
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
497+
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
498+
golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ=
499+
golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ=
500+
golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8=
501+
golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8=
502+
golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE=
461503
golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE=
462504
golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4=
463505
golang.org/x/time v0.15.0 h1:bbrp8t3bGUeFOx08pvsMYRTCVSMk89u4tKbNOZbp88U=
464506
golang.org/x/time v0.15.0/go.mod h1:Y4YMaQmXwGQZoFaVFk4YpCt4FLQMYKZe9oeV/f4MSno=
507+
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
508+
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
509+
golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc=
510+
golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU=
465511
golang.org/x/tools v0.45.0 h1:18qN3FAooORvApf5XjCXgsuayZOEtXf6JK18I3+ONa8=
466512
golang.org/x/tools v0.45.0/go.mod h1:LuUGqqaXcXMEFEruIVJVm5mgDD8vww/z/SR1gQ4uE/0=
513+
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
467514
golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
468515
gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4=
469516
gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E=

internal/outbox/kafka_publisher.go

Lines changed: 232 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,232 @@
1+
package outbox
2+
3+
import (
4+
"context"
5+
"encoding/json"
6+
"fmt"
7+
"os"
8+
"strings"
9+
"time"
10+
11+
"github.com/segmentio/kafka-go"
12+
)
13+
14+
type kafkaWriter interface {
15+
WriteMessages(ctx context.Context, msgs ...kafka.Message) error
16+
}
17+
18+
type kafkaWriterFactory func(brokers []string, topic string, ack kafka.RequiredAcks) kafkaWriter
19+
20+
type kafkaPublisherConfig struct {
21+
brokers []string
22+
topicMapping map[string]string
23+
defaultAck kafka.RequiredAcks
24+
topicAcks map[string]kafka.RequiredAcks
25+
}
26+
27+
// KafkaPublisher routes selected outbox topics to Kafka while preserving an
28+
// HTTP fallback for unconfigured topics.
29+
type KafkaPublisher struct {
30+
brokers []string
31+
topicMapping map[string]string
32+
topicAcks map[string]kafka.RequiredAcks
33+
defaultAck kafka.RequiredAcks
34+
writerFactory kafkaWriterFactory
35+
fallback Publisher
36+
}
37+
38+
func NewKafkaPublisher(brokers []string, topicMapping map[string]string, fallback Publisher) *KafkaPublisher {
39+
return &KafkaPublisher{
40+
brokers: append([]string(nil), brokers...),
41+
topicMapping: cloneTopicMapping(topicMapping),
42+
topicAcks: map[string]kafka.RequiredAcks{},
43+
defaultAck: kafka.RequireOne,
44+
writerFactory: newKafkaWriter,
45+
fallback: fallback,
46+
}
47+
}
48+
49+
func NewKafkaPublisherFromEnv(fallback Publisher) (*KafkaPublisher, error) {
50+
cfg, err := parseKafkaConfigFromEnv()
51+
if err != nil {
52+
return nil, err
53+
}
54+
55+
pub := NewKafkaPublisher(cfg.brokers, cfg.topicMapping, fallback)
56+
pub.defaultAck = cfg.defaultAck
57+
pub.topicAcks = cfg.topicAcks
58+
return pub, nil
59+
}
60+
61+
func (p *KafkaPublisher) Publish(ctx context.Context, event *Event) error {
62+
if event == nil {
63+
return fmt.Errorf("kafka publisher received nil event")
64+
}
65+
66+
topic, ok := p.resolveTopic(event)
67+
if !ok {
68+
if p.fallback != nil {
69+
return p.fallback.Publish(ctx, event)
70+
}
71+
return nil
72+
}
73+
74+
if len(p.brokers) == 0 {
75+
return fmt.Errorf("kafka publisher has no brokers configured")
76+
}
77+
78+
writer := p.writerFactory(p.brokers, topic, p.resolveAck(topic))
79+
if writer == nil {
80+
return fmt.Errorf("kafka publisher writer factory returned nil")
81+
}
82+
83+
payload := map[string]interface{}{
84+
"id": event.ID,
85+
"type": event.EventType,
86+
"data": event.EventData,
87+
"occurred_at": event.OccurredAt,
88+
"aggregate_id": safeString(event.AggregateID),
89+
"aggregate_type": safeString(event.AggregateType),
90+
"version": event.Version,
91+
}
92+
93+
body, err := json.Marshal(payload)
94+
if err != nil {
95+
return fmt.Errorf("failed to marshal kafka payload: %w", err)
96+
}
97+
98+
start := time.Now()
99+
err = writer.WriteMessages(ctx, kafka.Message{Key: []byte(event.ID.String()), Value: body})
100+
latency := time.Since(start).Seconds()
101+
if OutboxKafkaProduceLatency != nil {
102+
OutboxKafkaProduceLatency.WithLabelValues(topic).Observe(latency)
103+
}
104+
if err != nil {
105+
if OutboxKafkaErrorsTotal != nil {
106+
OutboxKafkaErrorsTotal.WithLabelValues(topic, "publish_error").Inc()
107+
}
108+
return fmt.Errorf("kafka publish failed for topic %s: %w", topic, err)
109+
}
110+
111+
return nil
112+
}
113+
114+
func (p *KafkaPublisher) resolveTopic(event *Event) (string, bool) {
115+
if event == nil {
116+
return "", false
117+
}
118+
if topic, ok := p.topicMapping[event.EventType]; ok && topic != "" {
119+
return topic, true
120+
}
121+
if len(event.EventData) == 0 {
122+
return "", false
123+
}
124+
var eventData EventData
125+
if err := json.Unmarshal(event.EventData, &eventData); err == nil && eventData.Type != "" {
126+
if topic, ok := p.topicMapping[eventData.Type]; ok && topic != "" {
127+
return topic, true
128+
}
129+
}
130+
return "", false
131+
}
132+
133+
func (p *KafkaPublisher) resolveAck(topic string) kafka.RequiredAcks {
134+
if ack, ok := p.topicAcks[topic]; ok {
135+
return ack
136+
}
137+
return p.defaultAck
138+
}
139+
140+
func newKafkaWriter(brokers []string, topic string, ack kafka.RequiredAcks) kafkaWriter {
141+
return &kafka.Writer{
142+
Addr: kafka.TCP(brokers...),
143+
Topic: topic,
144+
Balancer: &kafka.LeastBytes{},
145+
RequiredAcks: ack,
146+
Compression: kafka.Snappy,
147+
WriteTimeout: 10 * time.Second,
148+
ReadTimeout: 10 * time.Second,
149+
AllowAutoTopicCreation: true,
150+
}
151+
}
152+
153+
func parseKafkaConfigFromEnv() (kafkaPublisherConfig, error) {
154+
cfg := kafkaPublisherConfig{
155+
defaultAck: kafka.RequireOne,
156+
topicAcks: map[string]kafka.RequiredAcks{},
157+
}
158+
159+
for _, broker := range strings.Split(os.Getenv("OUTBOX_KAFKA_BROKERS"), ",") {
160+
broker = strings.TrimSpace(broker)
161+
if broker != "" {
162+
cfg.brokers = append(cfg.brokers, broker)
163+
}
164+
}
165+
166+
if topicMap := os.Getenv("OUTBOX_KAFKA_TOPIC_MAP"); topicMap != "" {
167+
cfg.topicMapping = map[string]string{}
168+
for _, entry := range strings.Split(topicMap, ",") {
169+
entry = strings.TrimSpace(entry)
170+
if entry == "" {
171+
continue
172+
}
173+
parts := strings.SplitN(entry, "=", 2)
174+
if len(parts) != 2 {
175+
return kafkaPublisherConfig{}, fmt.Errorf("invalid OUTBOX_KAFKA_TOPIC_MAP entry %q", entry)
176+
}
177+
cfg.topicMapping[strings.TrimSpace(parts[0])] = strings.TrimSpace(parts[1])
178+
}
179+
}
180+
181+
if ackValue := strings.TrimSpace(os.Getenv("OUTBOX_KAFKA_ACKS")); ackValue != "" {
182+
ack, err := parseKafkaAck(ackValue)
183+
if err != nil {
184+
return kafkaPublisherConfig{}, err
185+
}
186+
cfg.defaultAck = ack
187+
}
188+
189+
if ackValue := strings.TrimSpace(os.Getenv("OUTBOX_KAFKA_TOPIC_ACKS")); ackValue != "" {
190+
for _, entry := range strings.Split(ackValue, ",") {
191+
entry = strings.TrimSpace(entry)
192+
if entry == "" {
193+
continue
194+
}
195+
parts := strings.SplitN(entry, "=", 2)
196+
if len(parts) != 2 {
197+
return kafkaPublisherConfig{}, fmt.Errorf("invalid OUTBOX_KAFKA_TOPIC_ACKS entry %q", entry)
198+
}
199+
ack, err := parseKafkaAck(strings.TrimSpace(parts[1]))
200+
if err != nil {
201+
return kafkaPublisherConfig{}, err
202+
}
203+
cfg.topicAcks[strings.TrimSpace(parts[0])] = ack
204+
}
205+
}
206+
207+
return cfg, nil
208+
}
209+
210+
func parseKafkaAck(value string) (kafka.RequiredAcks, error) {
211+
switch strings.ToLower(strings.TrimSpace(value)) {
212+
case "0", "none", "no_ack":
213+
return kafka.RequireNone, nil
214+
case "1", "one", "requireone":
215+
return kafka.RequireOne, nil
216+
case "all", "requireall", "-1":
217+
return kafka.RequireAll, nil
218+
default:
219+
return kafka.RequireOne, fmt.Errorf("unsupported kafka ack value %q", value)
220+
}
221+
}
222+
223+
func cloneTopicMapping(mapping map[string]string) map[string]string {
224+
if len(mapping) == 0 {
225+
return map[string]string{}
226+
}
227+
cloned := make(map[string]string, len(mapping))
228+
for key, value := range mapping {
229+
cloned[key] = value
230+
}
231+
return cloned
232+
}

0 commit comments

Comments
 (0)