Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -796,6 +796,7 @@ async def register_schema_full_response(
# The registered schema may not be fully populated
s = registered_schema.schema if registered_schema.schema.schema_str is not None else schema
self._cache.set_schema(subject_name, registered_schema.schema_id, registered_schema.guid, s)
await self.clear_latest_caches(subject=subject_name)

return registered_schema

Expand Down Expand Up @@ -1613,10 +1614,17 @@ async def get_contexts(self, offset: int = 0, limit: int = -1) -> List[str]:
result = await self._rest_client.get('contexts', query={'offset': offset, 'limit': limit})
return result

async def clear_latest_caches(self):
async def clear_latest_caches(self, subject: Optional[str] = None) -> None:
"""Clear latest-version and latest-with-metadata caches, optionally for one subject."""
async with self._latest_lock:
self._latest_version_cache.clear()
self._latest_with_metadata_cache.clear()
if subject is None:
self._latest_version_cache.clear()
self._latest_with_metadata_cache.clear()
else:
self._latest_version_cache.pop(subject, None)
for cache_key in list(self._latest_with_metadata_cache):
if cache_key[0] == subject:
del self._latest_with_metadata_cache[cache_key]

async def clear_caches(self):
async with self._latest_lock:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -793,6 +793,7 @@ def register_schema_full_response(
# The registered schema may not be fully populated
s = registered_schema.schema if registered_schema.schema.schema_str is not None else schema
self._cache.set_schema(subject_name, registered_schema.schema_id, registered_schema.guid, s)
self.clear_latest_caches(subject=subject_name)

return registered_schema

Expand Down Expand Up @@ -1600,10 +1601,17 @@ def get_contexts(self, offset: int = 0, limit: int = -1) -> List[str]:
result = self._rest_client.get('contexts', query={'offset': offset, 'limit': limit})
return result

def clear_latest_caches(self):
def clear_latest_caches(self, subject: Optional[str] = None) -> None:
"""Clear latest-version and latest-with-metadata caches, optionally for one subject."""
with self._latest_lock:
self._latest_version_cache.clear()
self._latest_with_metadata_cache.clear()
if subject is None:
self._latest_version_cache.clear()
self._latest_with_metadata_cache.clear()
else:
self._latest_version_cache.pop(subject, None)
for cache_key in list(self._latest_with_metadata_cache):
if cache_key[0] == subject:
del self._latest_with_metadata_cache[cache_key]

def clear_caches(self):
with self._latest_lock:
Expand Down
56 changes: 55 additions & 1 deletion tests/schema_registry/_async/test_api_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@

import pytest

from confluent_kafka.schema_registry.common.schema_registry_client import SchemaVersion
from confluent_kafka.schema_registry.common.schema_registry_client import RegisteredSchema, SchemaVersion
from confluent_kafka.schema_registry.error import SchemaRegistryError
from confluent_kafka.schema_registry.schema_registry_client import AsyncSchemaRegistryClient, Schema
from tests.schema_registry.conftest import COUNTER, SCHEMA, SCHEMA_ID, SUBJECTS, USERINFO, VERSION, VERSIONS
Expand Down Expand Up @@ -94,6 +94,60 @@ async def test_register_schema_full_response_recall(mock_schema_registry, load_a
assert result.schema_id == SCHEMA_ID


async def test_register_schema_clears_latest_caches(mock_schema_registry, load_avsc):
sr = AsyncSchemaRegistryClient({'url': TEST_URL})
schema = Schema(load_avsc('basic_schema.avsc'), schema_type='AVRO')
cached = RegisteredSchema('test-key', VERSION - 1, SCHEMA_ID, None, schema)
sr._latest_version_cache['test-key'] = cached
sr._latest_with_metadata_cache[('test-key', frozenset({('owner', 'team')}), False)] = cached
other_cache_key = ('other-key', frozenset({('owner', 'other-team')}), False)
sr._latest_with_metadata_cache[other_cache_key] = cached

await sr.register_schema('test-key', schema)

assert 'test-key' not in sr._latest_version_cache
assert ('test-key', frozenset({('owner', 'team')}), False) not in sr._latest_with_metadata_cache
assert sr._latest_with_metadata_cache[other_cache_key] == cached


async def test_register_schema_cache_hit_keeps_latest_caches(mock_schema_registry, load_avsc):
sr = AsyncSchemaRegistryClient({'url': TEST_URL})
schema = Schema(load_avsc('basic_schema.avsc'), schema_type='AVRO')
cached = RegisteredSchema('test-key', VERSION, SCHEMA_ID, None, schema)
cache_key = ('test-key', frozenset({('owner', 'team')}), False)

await sr.register_schema('test-key', schema)
sr._latest_version_cache['test-key'] = cached
sr._latest_with_metadata_cache[cache_key] = cached

await sr.register_schema('test-key', schema)

assert sr._latest_version_cache['test-key'] == cached
assert sr._latest_with_metadata_cache[cache_key] == cached


async def test_clear_latest_caches(mock_schema_registry, load_avsc):
sr = AsyncSchemaRegistryClient({'url': TEST_URL})
schema = Schema(load_avsc('basic_schema.avsc'), schema_type='AVRO')
cached = RegisteredSchema('test-key', VERSION, SCHEMA_ID, None, schema)
sr._latest_version_cache['test-key'] = cached
sr._latest_version_cache['other-key'] = cached
sr._latest_with_metadata_cache[('test-key', frozenset(), False)] = cached
sr._latest_with_metadata_cache[('other-key', frozenset(), False)] = cached

await sr.clear_latest_caches(subject='test-key')

assert 'test-key' not in sr._latest_version_cache
assert ('test-key', frozenset(), False) not in sr._latest_with_metadata_cache
assert sr._latest_version_cache['other-key'] == cached
assert sr._latest_with_metadata_cache[('other-key', frozenset(), False)] == cached

await sr.clear_latest_caches()

assert not sr._latest_version_cache
assert not sr._latest_with_metadata_cache


async def test_register_schema_incompatible(mock_schema_registry, load_avsc):
conf = {'url': TEST_URL}
sr = AsyncSchemaRegistryClient(conf)
Expand Down
56 changes: 55 additions & 1 deletion tests/schema_registry/_sync/test_api_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@

import pytest

from confluent_kafka.schema_registry.common.schema_registry_client import SchemaVersion
from confluent_kafka.schema_registry.common.schema_registry_client import RegisteredSchema, SchemaVersion
from confluent_kafka.schema_registry.error import SchemaRegistryError
from confluent_kafka.schema_registry.schema_registry_client import Schema, SchemaRegistryClient
from tests.schema_registry.conftest import COUNTER, SCHEMA, SCHEMA_ID, SUBJECTS, USERINFO, VERSION, VERSIONS
Expand Down Expand Up @@ -94,6 +94,60 @@ def test_register_schema_full_response_recall(mock_schema_registry, load_avsc):
assert result.schema_id == SCHEMA_ID


def test_register_schema_clears_latest_caches(mock_schema_registry, load_avsc):
sr = SchemaRegistryClient({'url': TEST_URL})
schema = Schema(load_avsc('basic_schema.avsc'), schema_type='AVRO')
cached = RegisteredSchema('test-key', VERSION - 1, SCHEMA_ID, None, schema)
sr._latest_version_cache['test-key'] = cached
sr._latest_with_metadata_cache[('test-key', frozenset({('owner', 'team')}), False)] = cached
other_cache_key = ('other-key', frozenset({('owner', 'other-team')}), False)
sr._latest_with_metadata_cache[other_cache_key] = cached

sr.register_schema('test-key', schema)

assert 'test-key' not in sr._latest_version_cache
assert ('test-key', frozenset({('owner', 'team')}), False) not in sr._latest_with_metadata_cache
assert sr._latest_with_metadata_cache[other_cache_key] == cached


def test_register_schema_cache_hit_keeps_latest_caches(mock_schema_registry, load_avsc):
sr = SchemaRegistryClient({'url': TEST_URL})
schema = Schema(load_avsc('basic_schema.avsc'), schema_type='AVRO')
cached = RegisteredSchema('test-key', VERSION, SCHEMA_ID, None, schema)
cache_key = ('test-key', frozenset({('owner', 'team')}), False)

sr.register_schema('test-key', schema)
sr._latest_version_cache['test-key'] = cached
sr._latest_with_metadata_cache[cache_key] = cached

sr.register_schema('test-key', schema)

assert sr._latest_version_cache['test-key'] == cached
assert sr._latest_with_metadata_cache[cache_key] == cached


def test_clear_latest_caches(mock_schema_registry, load_avsc):
sr = SchemaRegistryClient({'url': TEST_URL})
schema = Schema(load_avsc('basic_schema.avsc'), schema_type='AVRO')
cached = RegisteredSchema('test-key', VERSION, SCHEMA_ID, None, schema)
sr._latest_version_cache['test-key'] = cached
sr._latest_version_cache['other-key'] = cached
sr._latest_with_metadata_cache[('test-key', frozenset(), False)] = cached
sr._latest_with_metadata_cache[('other-key', frozenset(), False)] = cached

sr.clear_latest_caches(subject='test-key')

assert 'test-key' not in sr._latest_version_cache
assert ('test-key', frozenset(), False) not in sr._latest_with_metadata_cache
assert sr._latest_version_cache['other-key'] == cached
assert sr._latest_with_metadata_cache[('other-key', frozenset(), False)] == cached

sr.clear_latest_caches()

assert not sr._latest_version_cache
assert not sr._latest_with_metadata_cache


def test_register_schema_incompatible(mock_schema_registry, load_avsc):
conf = {'url': TEST_URL}
sr = SchemaRegistryClient(conf)
Expand Down