Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,10 @@
*/
package org.apache.storm.kafka.bolt;

import org.apache.storm.kafka.mapper.FieldNameBasedTupleToKafkaMapper;
import org.apache.storm.kafka.mapper.TupleToKafkaMapper;
import org.apache.storm.kafka.selector.DefaultTopicSelector;
import org.apache.storm.kafka.selector.KafkaTopicSelector;
import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
Expand All @@ -29,10 +33,6 @@
import org.apache.kafka.clients.producer.Callback;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.storm.kafka.bolt.mapper.FieldNameBasedTupleToKafkaMapper;
import org.apache.storm.kafka.bolt.mapper.TupleToKafkaMapper;
import org.apache.storm.kafka.bolt.selector.DefaultTopicSelector;
import org.apache.storm.kafka.bolt.selector.KafkaTopicSelector;
import java.util.concurrent.Future;
import java.util.concurrent.ExecutionException;
import java.util.Map;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,9 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.storm.kafka.bolt.mapper;
package org.apache.storm.kafka.mapper;

import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.ITuple;

public class FieldNameBasedTupleToKafkaMapper<K,V> implements TupleToKafkaMapper<K, V> {

Expand All @@ -36,13 +36,13 @@ public FieldNameBasedTupleToKafkaMapper(String boltKeyField, String boltMessageF
}

@Override
public K getKeyFromTuple(Tuple tuple) {
public K getKeyFromTuple(ITuple tuple) {
//for backward compatibility, we return null when key is not present.
return tuple.contains(boltKeyField) ? (K) tuple.getValueByField(boltKeyField) : null;
}

@Override
public V getMessageFromTuple(Tuple tuple) {
public V getMessageFromTuple(ITuple tuple) {
return (V) tuple.getValueByField(boltMessageField);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,9 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.storm.kafka.bolt.mapper;
package org.apache.storm.kafka.mapper;

import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.ITuple;

import java.io.Serializable;

Expand All @@ -27,6 +27,6 @@
* @param <V> type of value.
*/
public interface TupleToKafkaMapper<K,V> extends Serializable {
K getKeyFromTuple(Tuple tuple);
V getMessageFromTuple(Tuple tuple);
K getKeyFromTuple(ITuple tuple);
V getMessageFromTuple(ITuple tuple);
}
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,9 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.storm.kafka.bolt.selector;
package org.apache.storm.kafka.selector;

import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.ITuple;

public class DefaultTopicSelector implements KafkaTopicSelector {

Expand All @@ -28,7 +28,7 @@ public DefaultTopicSelector(final String topicName) {
}

@Override
public String getTopic(Tuple tuple) {
public String getTopic(ITuple tuple) {
return topicName;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,9 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.storm.kafka.bolt.selector;
package org.apache.storm.kafka.selector;

import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.ITuple;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand All @@ -37,7 +37,7 @@ public FieldIndexTopicSelector(int fieldIndex, String defaultTopicName) {
}

@Override
public String getTopic(Tuple tuple) {
public String getTopic(ITuple tuple) {
if (fieldIndex < tuple.size()) {
return tuple.getString(fieldIndex);
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,9 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.storm.kafka.bolt.selector;
package org.apache.storm.kafka.selector;

import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.ITuple;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand All @@ -38,7 +38,7 @@ public FieldNameTopicSelector(String fieldName, String defaultTopicName) {
}

@Override
public String getTopic(Tuple tuple) {
public String getTopic(ITuple tuple) {
if (tuple.contains(fieldName)) {
return tuple.getStringByField(fieldName);
} else {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,12 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.storm.kafka.bolt.selector;
package org.apache.storm.kafka.selector;

import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.ITuple;

import java.io.Serializable;

public interface KafkaTopicSelector extends Serializable {
String getTopic(Tuple tuple);
String getTopic(ITuple tuple);
}
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@
*/
package org.apache.storm.kafka.trident;

import org.apache.storm.kafka.mapper.TupleToKafkaMapper;
import org.apache.storm.kafka.selector.KafkaTopicSelector;
import org.apache.storm.task.OutputCollector;
import org.apache.storm.topology.FailedException;
import org.apache.commons.lang.Validate;
Expand All @@ -25,8 +27,6 @@
import org.apache.kafka.clients.producer.RecordMetadata;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.storm.kafka.trident.mapper.TridentTupleToKafkaMapper;
import org.apache.storm.kafka.trident.selector.KafkaTopicSelector;
import org.apache.storm.trident.operation.TridentCollector;
import org.apache.storm.trident.state.State;
import org.apache.storm.trident.tuple.TridentTuple;
Expand All @@ -42,10 +42,10 @@ public class TridentKafkaState implements State {
private KafkaProducer producer;
private OutputCollector collector;

private TridentTupleToKafkaMapper mapper;
private TupleToKafkaMapper mapper;
private KafkaTopicSelector topicSelector;

public TridentKafkaState withTridentTupleToKafkaMapper(TridentTupleToKafkaMapper mapper) {
public TridentKafkaState withTridentTupleToKafkaMapper(TupleToKafkaMapper mapper) {
this.mapper = mapper;
return this;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,11 @@
*/
package org.apache.storm.kafka.trident;

import org.apache.storm.kafka.mapper.TupleToKafkaMapper;
import org.apache.storm.kafka.selector.KafkaTopicSelector;
import org.apache.storm.task.IMetricsContext;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.storm.kafka.trident.mapper.TridentTupleToKafkaMapper;
import org.apache.storm.kafka.trident.selector.KafkaTopicSelector;
import org.apache.storm.trident.state.State;
import org.apache.storm.trident.state.StateFactory;

Expand All @@ -32,11 +32,11 @@ public class TridentKafkaStateFactory implements StateFactory {

private static final Logger LOG = LoggerFactory.getLogger(TridentKafkaStateFactory.class);

private TridentTupleToKafkaMapper mapper;
private TupleToKafkaMapper mapper;
private KafkaTopicSelector topicSelector;
private Properties producerProperties = new Properties();

public TridentKafkaStateFactory withTridentTupleToKafkaMapper(TridentTupleToKafkaMapper mapper) {
public TridentKafkaStateFactory withTridentTupleToKafkaMapper(TupleToKafkaMapper mapper) {
this.mapper = mapper;
return this;
}
Expand Down

This file was deleted.

This file was deleted.

This file was deleted.

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -17,16 +17,16 @@
*/
package org.apache.storm.kafka;

import org.apache.storm.kafka.mapper.FieldNameBasedTupleToKafkaMapper;
import org.apache.storm.kafka.mapper.TupleToKafkaMapper;
import org.apache.storm.kafka.selector.DefaultTopicSelector;
import org.apache.storm.kafka.selector.KafkaTopicSelector;
import org.apache.storm.tuple.Fields;
import kafka.javaapi.consumer.SimpleConsumer;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.apache.storm.kafka.trident.TridentKafkaState;
import org.apache.storm.kafka.trident.mapper.FieldNameBasedTupleToKafkaMapper;
import org.apache.storm.kafka.trident.mapper.TridentTupleToKafkaMapper;
import org.apache.storm.kafka.trident.selector.DefaultTopicSelector;
import org.apache.storm.kafka.trident.selector.KafkaTopicSelector;
import org.apache.storm.trident.tuple.TridentTuple;
import org.apache.storm.trident.tuple.TridentTupleView;

Expand All @@ -42,7 +42,7 @@ public class TridentKafkaTest {
public void setup() {
broker = new KafkaTestBroker();
simpleConsumer = TestUtils.getKafkaConsumer(broker);
TridentTupleToKafkaMapper mapper = new FieldNameBasedTupleToKafkaMapper("key", "message");
TupleToKafkaMapper mapper = new FieldNameBasedTupleToKafkaMapper("key", "message");
KafkaTopicSelector topicSelector = new DefaultTopicSelector(TestUtils.TOPIC);
state = new TridentKafkaState()
.withKafkaTopicSelector(topicSelector)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,13 +20,12 @@
import org.apache.storm.Config;
import org.apache.storm.LocalCluster;
import org.apache.storm.generated.StormTopology;
import org.apache.storm.kafka.mapper.FieldNameBasedTupleToKafkaMapper;
import org.apache.storm.kafka.selector.DefaultTopicSelector;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Values;
import com.google.common.collect.ImmutableMap;
import org.apache.storm.kafka.trident.TridentKafkaStateFactory;
import org.apache.storm.kafka.trident.TridentKafkaUpdater;
import org.apache.storm.kafka.trident.mapper.FieldNameBasedTupleToKafkaMapper;
import org.apache.storm.kafka.trident.selector.DefaultTopicSelector;
import org.apache.storm.trident.Stream;
import org.apache.storm.trident.TridentTopology;
import org.apache.storm.trident.testing.FixedBatchSpout;
Expand Down