MarcusKainth opened a new issue, #4053:
URL: https://github.com/apache/iggy/issues/4053
### Bug description
A sink configured with `schema = "avro"`, `"proto"` or `"flatbuffer"`
receives
its messages as the wrong `Payload` variant.
The runtime tags each batch with the decoder's schema rather than the type
the
decoder returned. Those three decoders all return `Payload::Json` under the
config the runtime gives them, so the tag and the bytes disagree. The SDK
rebuilds the payload from the tag, so a sink is handed a `Payload::Avro`
holding
JSON.
What happens next depends on the sink. Sinks that only handle
`Payload::Json` —
ClickHouse, Delta, Quickwit, Redshift, Doris, Iceberg — log the message as an
unsupported payload type and drop it. Sinks that accept the raw variants —
S3,
HTTP, SurrealDB — base64-encode the JSON text and store it as though it were
Avro bytes.
Expected: the sink receives `Payload::Json`, because that is what the decoder
produced.
**Root cause**
```rust
// core/connectors/runtime/src/sink.rs:359-363
let messages_metadata = MessagesMetadata {
partition_id,
current_offset,
schema: decoder.schema(), // Schema::Avro
};
```
The bytes sent alongside come from `Payload::try_into_vec()`, which for the
`Payload::Json` the Avro decoder returned gives JSON
(`core/connectors/sdk/src/lib.rs:163-165`). The same tag is put on
`RawMessages` at `:731`. The SDK rebuilds from it (`sdk/src/sink.rs:199`,
then
`lib.rs:263`: `Schema::Avro => Payload::Avro(value)`).
No connector config avoids this. The runtime builds the Avro decoder with
`..AvroConfig::default()` (`sink.rs:532-536`), so `extract_as_json` cannot be
turned off from a connector file, and Proto and FlatBuffer get
`Schema::decoder()` defaults with no knob at all.
Sources are unaffected: a source plugin tags its own
`ProducedMessages.schema`
and the encoder takes the `Payload` directly (`source.rs:430-435`), with no
rebuild step.
**Which schemas are affected**
| `schema` | decoder returns | batch tagged | sink receives |
|---|---|---|---|
| `json`, `text`, `raw` | matching variant | matching | works |
| `avro` | `Payload::Json` (`decoders/avro.rs:169`) | `Avro` |
`Payload::Avro` |
| `flatbuffer` | `Payload::Json` (`decoders/flatbuffer.rs:109`) |
`FlatBuffer` | `Payload::FlatBuffer` |
| `proto` | `Payload::Json`, if the message is an `Any`
(`decoders/proto.rs:529`) | `Proto` | `Payload::Raw` |
The three that work are the ones whose decoder output matches its own
`schema()`. Avro and FlatBuffer default to `extract_as_json: true`; Proto
with
no schema configured falls through to `decode_as_any` (`proto.rs:264-265`),
which also returns JSON. Proto ends up as `Payload::Raw` rather than
`Payload::Proto` because `try_into_payload` re-decodes the bytes as a
`prost_types::Any` and falls back to `Raw` when that fails (`lib.rs:260`).
**Impact**
For the Json-only sinks the batch is lost without an error being raised. The
ClickHouse sink ends up with an empty body and returns `Ok(())`
(`sinks/clickhouse_sink/src/sink.rs:146-152`), and offsets are committed with
the poll request (`AutoCommit::When(AutoCommitWhen::PollingMessages)`,
`runtime/src/sink.rs:522`) — the SDK documents that mode as sending "the
commit
with the poll request itself, before your code sees the batch"
(`core/sdk/src/clients/consumer.rs:501`). So the offsets advance before the
sink
runs, and restarting the connector does not redeliver.
The batch is still counted as consumed. A run of 100 messages logs `Consumed
100
messages`, and the connector reports `Processed 0 messages.` at close.
Related to #3950 and PR #3954, which bind the sink's return status and let a
sink defer the offset commit. Neither covers this: the ClickHouse sink
returns
`Ok(())` on an empty body, so even with `offset_commit = after_consuming` the
batch would be committed as successfully written.
**Suggested fix**
Tag the batch with the payload's type rather than the decoder's. Changing
each
decoder's `schema()` would not be enough on its own: transforms run after
decoding and can change the type again, and the Proto decoder returns either
`Json` or `Raw` depending on the path it takes.
Adding a `Payload::schema()` and using it where `MessagesMetadata` and
`RawMessages` are built (`runtime/src/sink.rs:362` and `:731`) covers all
three
formats in one place. I can send a PR either way; say which you prefer.
**Compatibility**
The fix changes what three sinks write. S3, HTTP and SurrealDB match on
`Payload::Avro` and base64-encode it, so an avro stream reaches them today as
base64 of JSON text, stored as though it were Avro bytes. Nothing reports
this.
After the fix they receive `Payload::Json` and write a JSON document, so
anyone
consuming their current output would see it change shape. That output is
corrupt.
The Json-only sinks go from dropping the batch to writing it correctly.
Streams
configured `json`, `text` or `raw` are unaffected, because the decoder's
schema
and the payload's type already agree and the tag does not move.
Two related gaps:
- `StreamConsumerConfig` only exposes `avro_schema_json` and
`avro_schema_path`,
so flatbuffer and proto cannot be given a schema through the runtime at
all.
- Nothing exercises a non-JSON schema through the runtime's sink path, and
all
ten configs under `runtime/example_config/connectors/` use `schema =
"json"`.
A test across the decoders would catch this class of problem.
### Affected area / component
Connectors
### Deployment
Compiled from source
### Versions
master @ `412014a` — `iggy-connectors` 0.5.0-edge.6, `iggy_connector_sdk`
0.4.0-edge.3, `iggy_connector_clickhouse_sink` 0.2.0-edge.4
### Hardware / environment
macOS arm64, local build; not hardware-dependent.
### Sample code
```toml
type = "sink"
key = "clickhouse"
enabled = true
version = 0
name = "ClickHouse sink"
path = "target/release/libiggy_connector_clickhouse_sink"
verbose = true
[[streams]]
stream = "repro"
topics = ["avro"]
schema = "avro"
avro_schema_json =
'{"type":"record","name":"Event","fields":[{"name":"id","type":"long"},{"name":"name","type":"string"}]}'
batch_length = 100
poll_interval = "5ms"
consumer_group = "repro-sink"
[plugin_config]
url = "http://localhost:8123"
database = "default"
table = "repro"
insert_format = "row_binary"
```
Producer — an ordinary `IggyProducer`, with each payload an Avro datum for
that
schema:
```rust
let schema = apache_avro::Schema::parse_str(SCHEMA)?;
let record = apache_avro::types::Value::Record(vec![
("id".to_owned(), apache_avro::types::Value::Long(id)),
("name".to_owned(),
apache_avro::types::Value::String(format!("row-{id}"))),
]);
let datum = apache_avro::writer::datum::GenericDatumWriter::builder(&schema)
.build()?
.write_value_to_vec(record)?;
let message = IggyMessage::builder().payload(datum.into()).build()?;
```
### Logs
100 Avro messages into a single-partition topic, ClickHouse sink in
`row_binary` mode:
```
INFO iggy_connectors::sink: Processing 100 messages for sink connector with
ID: 1
ERROR connector: connector_target="iggy_connector_clickhouse_sink::body"
RowBinary mode: skipping unsupported payload type at offset 0
ERROR ... at offset 1, 2, ... 99
ERROR connector: connector_target="iggy_connector_clickhouse_sink::sink"
ClickHouse sink ID: 1 — no serialisable messages in batch of 100
INFO iggy_connectors::sink: Consumed 100 messages in 675µs for sink
connector with ID: 1
INFO connector: ClickHouse sink ID: 1 closed. Processed 0 messages.
```
`SELECT count() FROM default.repro` returns 0.
Restarting the connector against the same consumer group redelivers nothing:
no
batch, no errors, still no rows. The offsets were committed at poll time, so
the
100 messages cannot be recovered without resetting them by hand.
Tagging the batch from the payload's type rather than the decoder's makes the
same configuration land all 100 rows with the expected values, which confirms
the cause.
### Iggy server config
Default server config; not involved.
### Reproduction
1. Run an Iggy server and a ClickHouse instance with
`CREATE TABLE repro (id Int64, name String) ENGINE = MergeTree ORDER BY
id`.
2. Produce Avro-encoded records matching the schema above to `repro/avro`.
3. Run `iggy-connectors` with the sink config above.
4. The log fills with `skipping unsupported payload type` and the table stays
empty.
5. Restart the connector: nothing is redelivered, and the table is still
empty.
### Contribution
- [x] I'm willing to submit a pull request to fix this bug
### Good first issue
- [ ] I think this could be a good first issue for a new contributor
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]