systemtest/producer_consumer.go (27 lines of code) (raw):
// Licensed to Elasticsearch B.V. under one or more contributor
// license agreements. See the NOTICE file distributed with
// this work for additional information regarding copyright
// ownership. Elasticsearch B.V. licenses this file to you under
// the Apache License, Version 2.0 (the "License"); you may
// not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
package systemtest
import (
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/elastic/apm-queue/v2/kafka"
)
func newKafkaProducer(t testing.TB, cfg kafka.ProducerConfig) *kafka.Producer {
cfg.CommonConfig = KafkaCommonConfig(t, cfg.CommonConfig)
producer, err := kafka.NewProducer(cfg)
require.NoError(t, err)
t.Cleanup(func() {
err := producer.Close()
assert.NoError(t, err)
})
return producer
}
func newKafkaConsumer(t testing.TB, cfg kafka.ConsumerConfig) *kafka.Consumer {
cfg.CommonConfig = KafkaCommonConfig(t, cfg.CommonConfig)
consumer, err := kafka.NewConsumer(cfg)
require.NoError(t, err)
t.Cleanup(func() {
err := consumer.Close()
assert.NoError(t, err)
})
return consumer
}