【发布时间】:2021-12-19 13:01:35
【问题描述】:
由于 RocksDB 仍然不支持 Apple Silicon,目前只能使用通过 Rosetta 的 x86_64 JDK,它比原生 JDK 慢 5 倍。 因此,我想用内存中的键值存储替换 RocksDB。 如何将 Kafka 配置为默认使用这样的内存存储?
【问题讨论】:
标签: apache-kafka apache-kafka-streams spring-kafka rocksdb
由于 RocksDB 仍然不支持 Apple Silicon,目前只能使用通过 Rosetta 的 x86_64 JDK,它比原生 JDK 慢 5 倍。 因此,我想用内存中的键值存储替换 RocksDB。 如何将 Kafka 配置为默认使用这样的内存存储?
【问题讨论】:
标签: apache-kafka apache-kafka-streams spring-kafka rocksdb
在将storeBuilder 添加到StreamsBuilder 时,您可以选择构建持久化 (rocksdb) 或 in-mem 存储。
final var storeBuilder = Stores.windowStoreBuilder(
// Stores.persistentWindowStore(storeName, Duration.ofMinutes(10), Duration.ofMinutes(1),
// false),
Stores.inMemoryWindowStore(storeName, Duration.ofMinutes(10), Duration.ofMinutes(1),
false),
Serdes.String(),
Serdes.String()
);
builder.addStateStore(storeBuilder);
确保覆盖测试用例上的storeBuilder。
【讨论】:
这与 jego 的回答类似,但我与供应商合作。 这是我配置它们的方式:
@Profile("prod || stage || test")
@Configuration
class PersistentStoreConfiguration {
@Bean
fun projektanhangStoreSupplier(): KeyValueBytesStoreSupplier = Stores.persistentKeyValueStore(ProjektanhangStore.NAME)
}
@Profile( "it || dev")
@Configuration
class ProjektInMemoryStoreConfiguration {
@Bean
fun projektanhangStoreSupplier(): KeyValueBytesStoreSupplier = Stores.inMemoryKeyValueStore(ProjektanhangStore.NAME)
}
这就是根据弹簧轮廓选择的供应商将在何处以及如何注入和使用。注意@Bean 和@Configuration 类名。
@Configuration
class ProjektAnhangStreamConfiguration {
@Inject
private lateinit var projektanhangStoreSupplier: KeyValueBytesStoreSupplier
@Bean
fun projektanhaenge() = Consumer<KStream<String, AnhangEvent>> {
it.map { _, v -> KeyValue(v.anhang.projektId, v) }
.groupByKey(Grouped.with(Serdes.StringSerde(), JsonSerde(AnhangEvent::class.java)))
.aggregate(
{ ProjektanhangAggregator() },
{ _, anhangEvent, aggregator ->
when (anhangEvent.action) {
CREATE -> aggregator.add(anhangEvent.anhang)
DELETE -> aggregator.remove(anhangEvent.anhang)
UPDATE -> aggregator.update(anhangEvent.anhang)
}
},
Materialized
.`as`<String, ProjektanhangAggregator>(projektanhangStoreSupplier)
.withKeySerde(Serdes.String())
.withValueSerde(JsonSerde(ProjektanhangAggregator::class.java))
)
}
}
【讨论】: