|
19 | 19 | import io.opentelemetry.sdk.common.CompletableResultCode; |
20 | 20 | import io.opentelemetry.sdk.trace.data.SpanData; |
21 | 21 | import java.time.Duration; |
| 22 | +import java.util.Collection; |
22 | 23 | import java.util.List; |
23 | 24 | import java.util.UUID; |
24 | 25 | import java.util.concurrent.TimeUnit; |
|
29 | 30 | import org.apache.kafka.clients.consumer.ConsumerRecord; |
30 | 31 | import org.apache.kafka.clients.consumer.ConsumerRecords; |
31 | 32 | import org.apache.kafka.clients.consumer.KafkaConsumer; |
| 33 | +import org.apache.kafka.clients.producer.MockProducer; |
32 | 34 | import org.apache.kafka.clients.producer.ProducerConfig; |
| 35 | +import org.apache.kafka.common.KafkaException; |
33 | 36 | import org.apache.kafka.common.errors.ApiException; |
34 | 37 | import org.apache.kafka.common.serialization.StringDeserializer; |
35 | 38 | import org.apache.kafka.common.serialization.StringSerializer; |
|
46 | 49 | @TestInstance(TestInstance.Lifecycle.PER_CLASS) |
47 | 50 | class KafkaSpanExporterIntegrationTest { |
48 | 51 | private static final DockerImageName KAFKA_TEST_IMAGE = |
49 | | - DockerImageName.parse("apache/kafka:3.8.1"); |
| 52 | + DockerImageName.parse("apache/kafka:3.9.1"); |
50 | 53 | private static final String TOPIC = "span_topic"; |
51 | 54 | private KafkaContainer kafka; |
52 | 55 | private KafkaConsumer<String, ExportTraceServiceRequest> consumer; |
@@ -155,6 +158,28 @@ void exportWhenProducerInError() { |
155 | 158 | testSubject.shutdown(); |
156 | 159 | } |
157 | 160 |
|
| 161 | + @Test |
| 162 | + void exportWhenProducerFailsToSend() { |
| 163 | + var mockProducer = new MockProducer<String, Collection<SpanData>>(); |
| 164 | + mockProducer.sendException = new KafkaException("Simulated kafka exception"); |
| 165 | + var testSubjectWithMockProducer = |
| 166 | + KafkaSpanExporter.newBuilder().setTopicName(TOPIC).setProducer(mockProducer).build(); |
| 167 | + |
| 168 | + ImmutableList<SpanData> spans = |
| 169 | + ImmutableList.of(makeBasicSpan("span-1"), makeBasicSpan("span-2")); |
| 170 | + |
| 171 | + CompletableResultCode actual = testSubjectWithMockProducer.export(spans); |
| 172 | + |
| 173 | + await() |
| 174 | + .untilAsserted( |
| 175 | + () -> { |
| 176 | + assertThat(actual.isSuccess()).isFalse(); |
| 177 | + assertThat(actual.isDone()).isTrue(); |
| 178 | + }); |
| 179 | + |
| 180 | + testSubjectWithMockProducer.shutdown(); |
| 181 | + } |
| 182 | + |
158 | 183 | @Test |
159 | 184 | void flush() { |
160 | 185 | CompletableResultCode actual = testSubject.flush(); |
|
0 commit comments