Share

Search Java Code Snippets


  Help us in improving the repository. Add new snippets through 'Submit Code Snippet ' link.





#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


Subscribe to Java News and Posts. Get latest updates and posts on Java from Buggybread.com
Enter your email address:
Delivered by FeedBurner
 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



Subscribe to Java News and Posts. Get latest updates and posts on Java from Buggybread.com
Enter your email address:
Delivered by FeedBurner



comments powered by Disqus