Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
76 commits
Select commit Hold shift + click to select a range
d06988c
First cut at cel functions
rayokota May 13, 2026
26292b4
Modify decimals.prec
rayokota May 14, 2026
2bdb119
Clean up dec fns
rayokota May 14, 2026
1228db0
Add decimal.sqrt
rayokota May 30, 2026
40902c3
Add decimal to double CEL function
rayokota Jun 3, 2026
ba8adf0
Add decimal mod CEL function
rayokota Jun 3, 2026
15b24b0
Add decimal min/max CEL functions
rayokota Jun 3, 2026
83a9636
Minor renaming
rayokota Jun 3, 2026
9cb25b2
Fix decimal ctor
rayokota Aug 16, 2026
9254635
Add decimal_util
rayokota Aug 19, 2026
c3d2ff7
Add cel variant fns
rayokota Aug 21, 2026
8079f3b
Add avro/proto variant
rayokota Aug 21, 2026
5e00baf
add variant builder
rayokota Aug 22, 2026
6f5fe63
Fix variant issues
rayokota Aug 22, 2026
7a0b977
Fix variant issues
rayokota Aug 22, 2026
706f595
Fix variant issues
rayokota Aug 23, 2026
b8e27f8
Fix decimal issues
rayokota Aug 23, 2026
f11e0b8
Fix variant issues
rayokota Aug 23, 2026
d88d3b2
Fix ts issues
rayokota Aug 23, 2026
168f165
Simplify ts ctor
rayokota Aug 23, 2026
8770c84
Add tests
rayokota Aug 24, 2026
ae75cc7
Minor fixes
rayokota Aug 25, 2026
ba58c9d
Fix variant null handling
rayokota Sep 5, 2026
a05c68e
Java parity for message transforms
rayokota Sep 6, 2026
40e458d
Java parity for message transforms
rayokota Sep 6, 2026
738dbe3
Fix null handling (c8)
rayokota Sep 7, 2026
7c26904
Fix nested types (c9)
rayokota Sep 7, 2026
4b87244
Fix nested types (c9)
rayokota Sep 7, 2026
9761d43
Fix overflow
rayokota Sep 7, 2026
dab4bb1
make style-fix
rayokota Sep 8, 2026
90857a0
Fix mypy
rayokota Sep 8, 2026
f7dc4d5
make style-fix
rayokota Sep 8, 2026
decb978
Minor fixes
rayokota Sep 8, 2026
9367d8c
Minor fix
rayokota Sep 8, 2026
f8d4d25
Minor fix
rayokota Sep 9, 2026
888d62d
Minor fix
rayokota Sep 9, 2026
5804997
Minor fix
rayokota Sep 9, 2026
1ccedc4
Minor fix
rayokota Sep 9, 2026
39c7112
Minor fixes
rayokota Sep 9, 2026
660f8ed
Minor fixes
rayokota Sep 9, 2026
ee9853c
Minor fixes
rayokota Sep 9, 2026
d2a9723
Minor fixes
rayokota Sep 9, 2026
2469c6e
Minor fixes
rayokota Sep 9, 2026
a14bec3
Minor fixes
rayokota Sep 9, 2026
ae44609
Minor fixes
rayokota Sep 9, 2026
85bc2f7
Fix decimal funcs
rayokota Sep 9, 2026
fadd0d9
Minor fixes
rayokota Sep 9, 2026
ee7172e
Decimal fixes
rayokota Sep 9, 2026
a8d6027
Move proto types
rayokota Sep 10, 2026
3829e5b
Add confluent/types alias
rayokota Sep 10, 2026
74d63b9
Minor fixes
rayokota Sep 10, 2026
470416c
Minor fixes
rayokota Sep 10, 2026
db504e5
Minor fixes
rayokota Sep 10, 2026
2205d0e
Fix variant.of ts
rayokota Sep 10, 2026
6590b17
Fix decimal funcs
rayokota Sep 10, 2026
5ef4ad2
Fix decimal funcs
rayokota Sep 10, 2026
a5d23a8
Minor fixes
rayokota Sep 10, 2026
d012a1d
Minor fixes
rayokota Sep 11, 2026
0886e03
Minor fixes
rayokota Sep 11, 2026
d88b6d2
Add import public to new confluent/type/Decimal
rayokota Sep 11, 2026
e5d824e
Add import public to new confluent/type/Decimal
rayokota Sep 11, 2026
4d5a065
make style-fix
rayokota Sep 11, 2026
22002b9
Fix mypy
rayokota Sep 11, 2026
b1fdc8c
Pass schema via field context if necessary
rayokota Sep 11, 2026
2d27cc9
make style-fix
rayokota Sep 12, 2026
c2f11fe
Minor fixes
rayokota Sep 12, 2026
c089cc6
Minor fixes
rayokota Sep 12, 2026
e33ba4a
make style-fix
rayokota Sep 13, 2026
15aa15e
Update CHANGELOG
rayokota Sep 15, 2026
bdad8bf
make style-fix
rayokota Sep 23, 2026
8d601f4
Fix flake8
rayokota Sep 24, 2026
ede5458
make style-fix
rayokota Sep 24, 2026
51f63e0
make style-fix
rayokota Sep 24, 2026
89a0cdc
Fix tests
rayokota Sep 24, 2026
2070739
Fix tests
rayokota Sep 24, 2026
147b61f
Fix tests
rayokota Sep 24, 2026
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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ v2.16.0 is a feature release with the following features, fixes and enhancements
`dlq.auto.flush=true` or give the serde its own `RuleRegistry` (closable on
shutdown) for durability.
- Add support for inline validation rules (#2326)
- Add Variant, Decimal, and Timestamp CEL functions (#2332)

### Fixes

Expand Down
11 changes: 11 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,11 @@ Homepage = "https://github.com/confluentinc/confluent-kafka-python"

[tool.mypy]
ignore_missing_imports = true
# The generated protobuf modules are not statically analysable: _builder injects the message
# classes into globals() at import time. Excluding them keeps them out of the build set, which
# is what lets follow_imports = "skip" below apply - it is ignored for files mypy was asked to
# check directly.
exclude = '_pb2\.py$'

[[tool.mypy.overrides]]
module = [
Expand All @@ -35,12 +40,18 @@ module = [
]
disable_error_code = ["assignment", "no-redef"]

# Generated protobuf modules build their message classes at import time, so the classes are
# not statically visible - neither in the generated file itself nor to anything annotating
# against them. follow_imports = "skip" makes the module Any, which covers both.
[[tool.mypy.overrides]]
module = [
"confluent_kafka.schema_registry.confluent.meta_pb2",
"confluent_kafka.schema_registry.confluent.type.decimal_pb2",
"confluent_kafka.schema_registry.confluent.type.variant_pb2",
"confluent_kafka.schema_registry.confluent.types.decimal_pb2",
]
ignore_errors = true
follow_imports = "skip"

[tool.black]
line-length = 120
Expand Down
2 changes: 1 addition & 1 deletion src/confluent_kafka/admin/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
from .._model import ElectionType as _ElectionType
from .._model import TopicCollection as _TopicCollection
from ..cimpl import KafkaException # noqa: F401
from ..cimpl import _AdminClientImpl # noqa: F401
from ..cimpl import ( # noqa: F401
CONFIG_SOURCE_DEFAULT_CONFIG,
CONFIG_SOURCE_DYNAMIC_BROKER_CONFIG,
Expand All @@ -53,7 +54,6 @@
NewTopic,
)
from ..cimpl import TopicPartition as _TopicPartition
from ..cimpl import _AdminClientImpl
from ._acl import AclOperation # noqa: F401
from ._acl import AclBinding, AclBindingFilter, AclPermissionType # noqa: F401
from ._cluster import DescribeClusterResult # noqa: F401
Expand Down
55 changes: 53 additions & 2 deletions src/confluent_kafka/schema_registry/common/avro.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
from fastavro import repository, validate
from fastavro.schema import load_schema

from confluent_kafka.schema_registry.confluent.type.variant_utils import Variant
from confluent_kafka.schema_registry.serde import (
VALIDATION_RULES_PROP,
FieldTransform,
Expand All @@ -25,6 +26,33 @@

from .schema_registry_client import RuleKind, Schema


# The Avro `variant` logical type: a record {metadata: bytes, value: bytes} carrying a
# Spark/Parquet Variant. A field with this logical type decodes to / encodes from a Variant,
# so serde consumers and CEL rules see a first-class Variant rather than raw bytes. fastavro
# keys logical handlers by "<avro-type>-<logicalType>" = "record-variant"; this is the Python
# counterpart of Java's io.confluent.avro.type.VariantConversion.
def _variant_from_avro(data, writer_schema, reader_schema=None): # noqa: ARG001
return Variant(bytes(data["value"]), bytes(data["metadata"]))
Comment thread
rayokota marked this conversation as resolved.


def _variant_to_avro(data, schema): # noqa: ARG001
if isinstance(data, Variant):
# standalone_value_bytes, not .value: a navigated sub-variant's own value starts at
# its position, and .value is the whole shared buffer.
return {"metadata": data.metadata, "value": data.standalone_value_bytes()}
Comment thread
rayokota marked this conversation as resolved.
return data
Comment thread
rayokota marked this conversation as resolved.


try:
from fastavro import read as _fastavro_read
from fastavro import write as _fastavro_write

_fastavro_read.LOGICAL_READERS["record-variant"] = _variant_from_avro
_fastavro_write.LOGICAL_WRITERS["record-variant"] = _variant_to_avro
Comment thread
rayokota marked this conversation as resolved.
except Exception: # pragma: no cover - guards against a fastavro API shape change
pass

__all__ = [
'AvroMessage',
'AvroSchema',
Expand Down Expand Up @@ -113,10 +141,14 @@
return load_schema("$root", repo=repo)


def transform(

Check failure on line 144 in src/confluent_kafka/schema_registry/common/avro.py

View check run for this annotation

SonarQube-Confluent / SonarQube Code Analysis

Refactor this function to reduce its Cognitive Complexity from 36 to the 15 allowed.

[S3776] Cognitive Complexity of functions should not be too high See more on https://sonarqube.confluent.io/project/issues?id=confluent-kafka-python&pullRequest=2332&issues=cbcc0d95-72d2-4938-8d7c-7f3c1d73bae9&open=cbcc0d95-72d2-4938-8d7c-7f3c1d73bae9
ctx: RuleContext, schema: AvroSchema, message: AvroMessage, field_transform: FieldTransform
) -> AvroMessage:
if message is None or schema is None:
# Only the schema being absent stops the walk. A `None` *value* is the null branch of a
# `["null", T]` union and has to reach the rule: the reference binds it as CEL null so a
# rule can guard with `value == null`, and returning early here skipped the rule entirely -
# indistinguishable, to the caller, from a rule that ran and passed.
if schema is None:
return message
field_ctx = ctx.current_field()
if field_ctx is not None:
Expand All @@ -126,7 +158,12 @@
if subschema is None:
return message
submessage = transform(ctx, subschema, submessage, field_transform)
if isinstance(message, tuple) and len(message) == 2:
# The branch a transformed value belongs to follows from the value, not from the branch
# it arrived on - the reference keeps no branch at all and resolves it from the datum.
# Keep the notation while its branch still accepts the result, so two same-shaped
# records are never swapped; drop it otherwise and let fastavro resolve, since
# ("null", x) would be written as null with x silently dropped.
if isinstance(message, tuple) and len(message) == 2 and _branch_accepts(subschema, submessage):
return (message[0], submessage)
return submessage
elif isinstance(schema, dict):
Expand All @@ -142,6 +179,11 @@
return message
return {key: transform(ctx, schema["values"], value, field_transform) for key, value in message.items()}
elif schema_type == 'record':
# A null record has no fields to walk. Guarded before the isinstance check below
# so a legitimate null does not log an "incompatible message type" warning; the
# reference guards the record case, and only the record case, the same way.
if message is None:
return message
if not isinstance(message, dict):
log.warning("Incompatible message type for record schema")
return message
Expand Down Expand Up @@ -364,6 +406,15 @@
return '.' not in name and not subschema.get("namespace") and branch_name.rsplit('.', 1)[-1] == name


def _branch_accepts(subschema: AvroSchema, message: AvroMessage) -> bool:
"""Whether ``subschema`` can hold ``message``, by the same test ``_resolve_union`` uses."""
try:
validate(message, _collapse_schema(deepcopy(subschema)))
return True
except: # noqa: E722

Check failure on line 414 in src/confluent_kafka/schema_registry/common/avro.py

View check run for this annotation

SonarQube-Confluent / SonarQube Code Analysis

Specify an exception class to catch or reraise the exception

[S5754] "SystemExit" should be re-raised See more on https://sonarqube.confluent.io/project/issues?id=confluent-kafka-python&pullRequest=2332&issues=f83da161-a828-4700-9d37-1ed0d11a5b58&open=f83da161-a828-4700-9d37-1ed0d11a5b58
return False


def _resolve_union(schema: AvroSchema, message: AvroMessage) -> Tuple[Optional[AvroSchema], AvroMessage]:
is_wrapped_union = isinstance(message, tuple) and len(message) == 2
is_typed_union = isinstance(message, dict) and '-type' in message
Expand Down
Loading