Aiven-Open / Aiven-Open/cloud-storage-connectors-for-apache-kafka

S3 source failing to serialise S3 sink files

Abierto
#512 6 comentarios 0 reacciones 0 asignados Ver en GitHub
bug conflicts
Lenguaje dominante
Java
Estrellas
58
Forks
39
Merge medio
2 d 10 h
PR fusionados (30 d)
5

Descripción

### Introduction

@Ahmedie-m opened this ticket with the information below. After several session Ahmedie determined that the issue is that S3 sink writes with GZIP by default while S3 source defaults to none.

The resolution of this ticket will be that all source and sink pairs will be evaluated to show that data written with the default values for the sink are read by the default values for the source.

### Initial report

I have a setup of Aiven's S3 source & S3 sink with the following configurations:

S3 Sink:
```
"file.name.template" = "{{topic}}-{{partition}}-{{start_offset}}"
"format.output.fields" = "value,key,offset,timestamp"
"format.output.fields.value.encoding" = "none"
"format.output.type" = "jsonl"
"key.converter" = "org.apache.kafka.connect.storage.StringConverter"
"value.converter" = "org.apache.kafka.connect.json.JsonConverter"
"value.converter.schemas.enable" = "false"
"key.converter.schemas.enable" = "false"
```

S3 Source:
```
"input.format" = "jsonl"
"file.name.template" = "{{topic}}-{{partition}}-{{start_offset}}"
"key.converter" = "org.apache.kafka.connect.storage.StringConverter"
"value.converter" = "org.apache.kafka.connect.json.JsonConverter"
"key.converter.schemas.enable" = "false"
"value.converter.schemas.enable" = "false"
```

The files are correctly being sent to the S3 bucket in `jsonl` format where each line is a new json object.

However when the S3 source attempts to serialize those files it throws the following:

```
[Worker-0f5b346ed0b78484a] [2025-08-20 12:14:30,963] ERROR [s3-to-kafka|task-0] Error trying to advance data: Converting byte[] to Kafka Connect data failed due to serialization error: (io.aiven.kafka.connect.common.source.input.JsonTransformer:164)
[Worker-0f5b346ed0b78484a] org.apache.kafka.connect.errors.DataException: Converting byte[] to Kafka Connect data failed due to serialization error:
[Worker-0f5b346ed0b78484a]   at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:333)
[Worker-0f5b346ed0b78484a]   at io.aiven.kafka.connect.common.source.input.JsonTransformer$1.doAdvance(JsonTransformer.java:88)
[Worker-0f5b346ed0b78484a]   at io.aiven.kafka.connect.common.source.input.Transformer$StreamSpliterator.tryAdvance(Transformer.java:162)
[Worker-0f5b346ed0b78484a]   at java.base/java.util.stream.StreamSpliterators$WrappingSpliterator.lambda$initPartialTraversalState$0(StreamSpliterators.java:292)
[Worker-0f5b346ed0b78484a]   at java.base/java.util.stream.StreamSpliterators$AbstractWrappingSpliterator.fillBuffer(StreamSpliterators.java:206)
[Worker-0f5b346ed0b78484a]   at java.base/java.util.stream.StreamSpliterators$AbstractWrappingSpliterator.doAdvance(StreamSpliterators.java:161)
[Worker-0f5b346ed0b78484a]   at java.base/java.util.stream.StreamSpliterators$WrappingSpliterator.tryAdvance(StreamSpliterators.java:298)
[Worker-0f5b346ed0b78484a]   at java.base/java.util.Spliterators$1Adapter.hasNext(Spliterators.java:681)
[Worker-0f5b346ed0b78484a]   at io.aiven.kafka.connect.common.source.AbstractSourceRecordIterator.hasNext(AbstractSourceRecordIterator.java:200)
[Worker-0f5b346ed0b78484a]   at io.aiven.kafka.connect.s3.source.S3SourceTask$1.hasNext(S3SourceTask.java:83)
[Worker-0f5b346ed0b78484a]   at org.apache.commons.collections4.iterators.FilterIterator.setNextObject(FilterIterator.java:174)
[Worker-0f5b346ed0b78484a]   at org.apache.commons.collections4.iterators.FilterIterator.hasNext(FilterIterator.java:86)
[Worker-0f5b346ed0b78484a]   at io.aiven.kafka.connect.common.source.AbstractSourceTask.tryAdd(AbstractSourceTask.java:187)
[Worker-0f5b346ed0b78484a]   at io.aiven.kafka.connect.common.source.AbstractSourceTask$2.run(AbstractSourceTask.java:131)
[Worker-0f5b346ed0b78484a]   at java.base/java.lang.Thread.run(Thread.java:840)
[Worker-0f5b346ed0b78484a] Caused by: org.apache.kafka.common.errors.SerializationException: com.fasterxml.jackson.core.JsonParseException: Invalid UTF-8 start byte 0xbf
[Worker-0f5b346ed0b78484a] at [Source: (byte[])"���Y�����{>F]����u�b�}��(�"; line: 1, column: 3]
[Worker-0f5b346ed0b78484a]   at org.apache.kafka.connect.json.JsonDeserializer.deserialize(JsonDeserializer.java:69)
[Worker-0f5b346ed0b78484a]   at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:331)
[Worker-0f5b346ed0b78484a]   ... 14 more
[Worker-0f5b346ed0b78484a] Caused by: com.fasterxml.jackson.core.JsonParseException: Invalid UTF-8 start byte 0xbf
[Worker-0f5b346ed0b78484a] at [Source: (byte[])"���Y�����{>F]����u�b�}��(�"; line: 1, column: 3]
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:2337)
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:713)
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._reportInvalidInitial(UTF8StreamJsonParser.java:3607)
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._decodeCharForError(UTF8StreamJsonParser.java:3350)
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._reportInvalidToken(UTF8StreamJsonParser.java:3582)
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._handleUnexpectedValue(UTF8StreamJsonParser.java:2688)
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._nextTokenNotInObject(UTF8StreamJsonParser.java:870)
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.core.json.UTF8StreamJsonParser.nextToken(UTF8StreamJsonParser.java:762)
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.databind.ObjectMapper._readTreeAndClose(ObjectMapper.java:4622)
[Worker-0f5b346ed0b78484a]   at com.fasterxml.jackson.databind.ObjectMapper.readTree(ObjectMapper.java:3056)
[Worker-0f5b346ed0b78484a]   at org.apache.kafka.connect.json.JsonDeserializer.deserialize(JsonDeserializer.java:67)
[Worker-0f5b346ed0b78484a]   ... 15 more
[Worker-0f5b346ed0b78484a] [2025-08-20 12:14:30,966] INFO [s3-to-kafka|task-0] No records found in tryAdd call (io.aiven.kafka.connect.s3.source.S3SourceTask:196)
[Worker-0f5b346ed0b78484a] [2025-08-20 12:14:34,443] INFO [s3-to-kafka|task-0|offsets] Committed all records through last poll() (io.aiven.kafka.connect.s3.source.S3SourceTask:128)
```

What could be the reason why? I double checked that I am using the correct property of `input.format` instead of `input.type` as both are referenced in the `README.md`

Versions:
Aiven Connectors: v3.4.0
Kafka Connect: 3.7.x

Guía de contribución

Abrir la guía de contribución

Evaluación

Este issue todavía no se ha evaluado.

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.