Generate random data and produce to Kafka.
(
kafka_port: int,
sr_port: int,
topic: str,
debezium: bool,
column_dicts: list[dict[str, Any]],
num_rows: int,
rng_seed: int,
)
| 92 | |
| 93 | |
| 94 | def _kafka_chunk( |
| 95 | kafka_port: int, |
| 96 | sr_port: int, |
| 97 | topic: str, |
| 98 | debezium: bool, |
| 99 | column_dicts: list[dict[str, Any]], |
| 100 | num_rows: int, |
| 101 | rng_seed: int, |
| 102 | ) -> int: |
| 103 | """Generate random data and produce to Kafka.""" |
| 104 | rng = random.Random(rng_seed) |
| 105 | columns = [ |
| 106 | Column(c["name"], c["type"], c["nullable"], c["default"], c.get("data_shape")) |
| 107 | for c in column_dicts |
| 108 | ] |
| 109 | |
| 110 | producer, serializer, key_serializer, sctx, ksctx, col_names = get_kafka_objects( |
| 111 | topic, |
| 112 | tuple(columns), |
| 113 | debezium, |
| 114 | sr_port, |
| 115 | kafka_port, |
| 116 | ) |
| 117 | |
| 118 | now_ms = int(time.time() * 1000) |
| 119 | source_struct: dict[str, Any] | None = None |
| 120 | if debezium: |
| 121 | source_struct = { |
| 122 | "version": "0", |
| 123 | "connector": "mysql", |
| 124 | "name": "materialize-generator", |
| 125 | "ts_ms": now_ms, |
| 126 | "snapshot": None, |
| 127 | "db": "db", |
| 128 | "sequence": None, |
| 129 | "table": topic.split(".")[-1], |
| 130 | "server_id": 0, |
| 131 | "gtid": None, |
| 132 | "file": "binlog.000001", |
| 133 | "pos": 0, |
| 134 | "row": 0, |
| 135 | "thread": None, |
| 136 | "query": None, |
| 137 | } |
| 138 | |
| 139 | producer.poll(0) |
| 140 | for _ in range(num_rows): |
| 141 | row = [col.kafka_value(rng) for col in columns] |
| 142 | while True: |
| 143 | try: |
| 144 | if debezium: |
| 145 | after_value = dict(zip(col_names, row)) |
| 146 | envelope_value = { |
| 147 | "before": None, |
| 148 | "after": after_value, |
| 149 | "source": source_struct, |
| 150 | "op": "c", |
| 151 | "ts_ms": now_ms, |
nothing calls this directly
no test coverage detected