#Java - Code Snippets for '#Kafka' - 8 code snippet(s) found |
|
Sample 1. Initializing Executor Service within Kafka | |
|
int numberOfThreads = 10;
ExecutorService executor = Executors.newFixedThreadPool(numberOfThreads);
|
|
Like Feedback kafka kafka executor |
|
|
Sample 2. Code Sample / Example / Snippet of org.apache.storm.kafka.ZkHosts | |
|
private TransactionalTridentKafkaSpout createKafkaSpout() {
ZkHosts hosts = new ZkHosts(zkUrl);
TridentKafkaConfig config = new TridentKafkaConfig(hosts, "test");
config.scheme = new SchemeAsMultiScheme(new StringScheme());
config.startOffsetTime = kafka.api.OffsetRequest.LatestTime();
return new TransactionalTridentKafkaSpout(config);
}
|
|
Like Feedback org.apache.storm.kafka.ZkHosts |
|
|
Sample 3. Code Sample / Example / Snippet of org.apache.storm.kafka.trident.TridentKafkaConfig | |
|
private TransactionalTridentKafkaSpout createKafkaSpout() {
ZkHosts hosts = new ZkHosts(zkUrl);
TridentKafkaConfig config = new TridentKafkaConfig(hosts, "test");
config.scheme = new SchemeAsMultiScheme(new StringScheme());
config.startOffsetTime = kafka.api.OffsetRequest.LatestTime();
return new TransactionalTridentKafkaSpout(config);
}
|
|
Like Feedback org.apache.storm.kafka.trident.TridentKafkaConfig |
|
|
Sample 4. Code Sample / Example / Snippet of org.apache.storm.kafka.bolt.KafkaBolt | |
|
public StormTopology buildProducerTopology(Properties prop) {
TopologyBuilder builder = new TopologyBuilder();
builder.setSpout("spout", new RandomSentenceSpout(), 2);
KafkaBolt bolt = new KafkaBolt().withProducerProperties(prop)
.withTopicSelector(new DefaultTopicSelector("test"))
.withTupleToKafkaMapper(new FieldNameBasedTupleToKafkaMapper("key", "word"));
builder.setBolt("forwardToKafka", bolt, 1).shuffleGrouping("spout");
return builder.createTopology();
}
|
|
Like Feedback org.apache.storm.kafka.bolt.KafkaBolt |
|
|
|
Sample 5. Code Sample / Example / Snippet of org.apache.storm.kafka.trident.GlobalPartitionInformation | |
|
public static GlobalPartitionInformation buildPartitionInfo(int numPartitions, int brokerPort) {
GlobalPartitionInformation globalPartitionInformation = new GlobalPartitionInformation(TOPIC);
for (int i = 0; i < numPartitions; i++) {
globalPartitionInformation.addPartition(i, Broker.fromString("broker-" + i + " :" + brokerPort));
}
return globalPartitionInformation;
}
|
|
Like Feedback org.apache.storm.kafka.trident.GlobalPartitionInformation |
|
|
Sample 6. Code Sample / Example / Snippet of org.apache.storm.kafka.DynamicBrokersReader | |
|
public void testErrorLogsWhenConfigIsMissing() throws Exception {
String connectionString = server.getConnectString();
Map conf = new HashMap();
conf.put(Config.STORM_ZOOKEEPER_SESSION_TIMEOUT, 1000);
conf.put(Config.STORM_ZOOKEEPER_RETRY_TIMES, 4);
conf.put(Config.STORM_ZOOKEEPER_RETRY_INTERVAL, 5);
DynamicBrokersReader dynamicBrokersReader1 = new DynamicBrokersReader(conf, connectionString, masterPath, topic);
}
|
|
Like Feedback org.apache.storm.kafka.DynamicBrokersReader |
|
|
Sample 7. Code Sample / Example / Snippet of org.apache.storm.kafka.Partition | |
|
public void generateTuplesWithMessageAndMetadataScheme() {
String value = "value";
Partition mockPartition = Mockito.mock(Partition.class);
mockPartition.partition = 0;
long offset = 0L;
MessageMetadataSchemeAsMultiScheme scheme = new MessageMetadataSchemeAsMultiScheme(new StringMessageAndMetadataScheme());
createTopicAndSendMessage(null, value);
ByteBufferMessageSet messageAndOffsets = getLastMessage();
for (MessageAndOffset msg : messageAndOffsets) {
Iterable<List<Object>> lists = KafkaUtils.generateTuples(scheme, msg.message(), mockPartition, offset);
List<Object> values = lists.iterator().next();
assertEquals("Message is incorrect", value, values.get(0));
assertEquals("Partition is incorrect", mockPartition.partition, values.get(1));
assertEquals("Offset is incorrect", offset, values.get(2));
}
}
|
|
Like Feedback org.apache.storm.kafka.Partition |
|
|
Sample 8. Code Sample / Example / Snippet of org.apache.storm.kafka.PartitionManager.KafkaMessageId | |
|
public void ack(Object msgId) {
KafkaMessageId id = (KafkaMessageId) msgId;
PartitionManager m = _coordinator.getManager(id.partition);
if (m != null) {
m.ack(id.offset);
}
}
|
|
Like Feedback org.apache.storm.kafka.PartitionManager.KafkaMessageId |
|
|