KafkaPublisherSpecification(
_outer_config: T_BrokerConfig,
specification_config: T_SpecificationConfig,
)
Bases: PublisherSpecification[KafkaBrokerConfig, KafkaPublisherSpecificationConfig]
Source code in faststream/_internal/endpoint/publisher/specification.py
| def __init__(
self,
_outer_config: "T_BrokerConfig",
specification_config: "T_SpecificationConfig",
) -> None:
super().__init__(_outer_config, specification_config)
self.calls: list[AnyCallable] = []
|
config instance-attribute
config = specification_config
include_in_schema property
calls instance-attribute
calls: list[AnyCallable] = []
get_schema
Source code in faststream/confluent/publisher/specification.py
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48 | def get_schema(self) -> dict[str, PublisherSpec]:
payloads = self.get_payloads()
return {
self.name: PublisherSpec(
address=self.topic,
description=self.config.description_,
operation=Operation(
message=Message(
title=f"{self.name}:Message",
payload=resolve_payloads(payloads, "Publisher"),
),
bindings=None,
),
bindings=ChannelBinding(
kafka=kafka.ChannelBinding(
topic=self.topic,
partitions=None,
replicas=None,
),
),
),
}
|
add_call
add_call(call: AnyCallable) -> None
Source code in faststream/_internal/endpoint/publisher/specification.py
| def add_call(self, call: "AnyCallable") -> None:
self.calls.append(call)
|
get_payloads
get_payloads() -> list[tuple[dict[str, Any], str]]
Source code in faststream/_internal/endpoint/publisher/specification.py
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89 | def get_payloads(self) -> list[tuple[dict[str, Any], str]]:
payloads: list[tuple[dict[str, Any], str]] = []
if self.config.schema_:
body = get_model_schema(
call=create_model(
"",
__config__=get_config_base(),
response__=(self.config.schema_, ...),
),
prefix=f"{self.name}:Message",
)
if body: # pragma: no branch
payloads.append((body, ""))
else:
di_state = self._outer_config.fd_config
for call in self.calls:
call_model = build_call_model(
call,
dependency_provider=di_state.provider,
serializer_cls=di_state._serializer,
)
if call_model.serializer:
response_type = next(
iter(call_model.serializer.response_option.values()),
).field_type
else:
response_type = None
if response_type is not None and response_type is not Parameter.empty:
body = get_model_schema(
create_model(
"",
__config__=get_config_base(),
response__=(response_type, ...),
),
prefix=f"{self.name}:Message",
)
if body:
payloads.append((body, to_camelcase(unwrap(call).__name__)))
return payloads
|