Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,12 @@ The format is based on [Keep a Changelog](http://keepachangelog.com/) and this p
- upgrade template
- clarify Producer and Consumer task documentation (ports, entity serialization, dataset advanced-options reference) and fix minor typos in parameter descriptions

### Fixed

- Consumer and Producer tasks no longer convert non-ASCII characters to unicode escape
sequences when wrapping consumed Kafka messages into JSON, or when producing messages
from a JSON dataset or from entities

## [3.5.0] 2026-08-24

### Changed
Expand Down
4 changes: 2 additions & 2 deletions cmem_plugin_kafka/kafka_handlers.py
Original file line number Diff line number Diff line change
Expand Up @@ -181,7 +181,7 @@ def _split_data(self, data: Response) -> Generator[KafkaMessage]:
content = None
if not tombstone:
_content = _message.get("content")
content = json.dumps(_content) if _content else None
content = json.dumps(_content, ensure_ascii=False) if _content else None
yield KafkaMessage(key=key, value=content, headers=headers, tombstone=tombstone)

def _aggregate_data(self) -> Generator:
Expand Down Expand Up @@ -361,7 +361,7 @@ def _split_data(self, data: Entities) -> Generator[KafkaMessage]:
for i, path in enumerate(paths):
values[path.path] = list(entity.values[i])
result["entity"] = {"uri": entity.uri, "values": values}
kafka_payload = json.dumps(result, indent=4)
kafka_payload = json.dumps(result, indent=4, ensure_ascii=False)
yield KafkaMessage(key=None, value=kafka_payload)

def _aggregate_data(self) -> Entities:
Expand Down
2 changes: 1 addition & 1 deletion cmem_plugin_kafka/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -337,7 +337,7 @@ def get_message_with_json_wrapper(message: KafkaMessage) -> str:
msg_with_wrapper = {"message": {"key": message.key, "content": json.loads(message.value)}}
if message.headers:
msg_with_wrapper["message"]["headers"] = dict(message.headers.items())
return json.dumps(msg_with_wrapper, cls=BytesEncoder)
return json.dumps(msg_with_wrapper, cls=BytesEncoder, ensure_ascii=False)


def get_kafka_statistics(json_data: str) -> dict:
Expand Down
41 changes: 41 additions & 0 deletions tests/test_split_data_unicode.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
"""Regression tests for non-ASCII character handling in _split_data().

These exercise KafkaJSONDataHandler and KafkaEntitiesDataHandler directly,
without a Kafka broker or a Corporate Memory connection, since _split_data()
only transforms in-memory data into KafkaMessage objects.
"""

import json

import httpx
from cmem_plugin_base.dataintegration.entity import Entities, Entity, EntityPath, EntitySchema
from cmem_plugin_base.dataintegration.plugins import PluginLogger

from cmem_plugin_kafka.kafka_handlers import KafkaEntitiesDataHandler, KafkaJSONDataHandler


def test_json_data_handler_split_data_keeps_unicode_characters() -> None:
"""Test that non-ASCII characters from the source dataset are not escaped"""
payload = [{"message": {"key": "1", "content": {"city": "Köln", "name": "Müller"}}}]
data = httpx.Response(200, content=json.dumps(payload).encode("utf-8"))
handler = KafkaJSONDataHandler(context=None, plugin_logger=PluginLogger()) # type: ignore[arg-type]

messages = list(handler._split_data(data)) # noqa: SLF001

assert len(messages) == 1
assert "\\u00f6" not in messages[0].value
assert "\\u00fc" not in messages[0].value
assert json.loads(messages[0].value) == {"city": "Köln", "name": "Müller"}


def test_entities_data_handler_split_data_keeps_unicode_characters() -> None:
"""Test that non-ASCII characters in entity values are not escaped"""
schema = EntitySchema(type_uri="urn:x-test", paths=[EntityPath(path="city")])
entities = Entities(entities=[Entity(uri="urn:x-1", values=[["Köln"]])], schema=schema)
handler = KafkaEntitiesDataHandler(context=None, plugin_logger=PluginLogger()) # type: ignore[arg-type]

messages = list(handler._split_data(entities)) # noqa: SLF001

assert len(messages) == 1
assert "\\u00f6" not in messages[0].value
assert json.loads(messages[0].value)["entity"]["values"]["city"] == ["Köln"]
15 changes: 14 additions & 1 deletion tests/test_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,20 @@

import json

from cmem_plugin_kafka.utils import get_kafka_statistics
from cmem_plugin_kafka.utils import (
KafkaMessage,
get_kafka_statistics,
get_message_with_json_wrapper,
)


def test_get_message_with_json_wrapper_keeps_unicode_characters() -> None:
"""Test that non-ASCII characters in the message value are not escaped"""
message = KafkaMessage(key="1", value='{"city": "Köln", "name": "Müller"}')
wrapped = get_message_with_json_wrapper(message)
assert "\\u00f6" not in wrapped
assert "\\u00fc" not in wrapped
assert json.loads(wrapped)["message"]["content"] == {"city": "Köln", "name": "Müller"}


def test_get_kafka_statistics() -> None:
Expand Down
Loading