Skip to content

Commit

Permalink
add opaque partitioned transactional spout to builder
Browse files Browse the repository at this point in the history
  • Loading branch information
Nathan Marz committed Feb 11, 2012
1 parent 693571f commit 3df25ab
Show file tree
Hide file tree
Showing 2 changed files with 15 additions and 4 deletions.
Expand Up @@ -19,7 +19,9 @@
import backtype.storm.topology.InputDeclarer;
import backtype.storm.topology.SpoutDeclarer;
import backtype.storm.topology.TopologyBuilder;
import backtype.storm.transactional.partitioned.IOpaquePartitionedTransactionalSpout;
import backtype.storm.transactional.partitioned.IPartitionedTransactionalSpout;
import backtype.storm.transactional.partitioned.OpaquePartitionedTransactionalSpoutExecutor;
import backtype.storm.transactional.partitioned.PartitionedTransactionalSpoutExecutor;
import backtype.storm.tuple.Fields;
import java.util.ArrayList;
Expand Down Expand Up @@ -53,17 +55,26 @@ public TransactionalTopologyBuilder(String id, String spoutId, ITransactionalSpo
_spoutParallelism = spoutParallelism;
}

public TransactionalTopologyBuilder(String id, String spoutId, ITransactionalSpout spout) {
this(id, spoutId, spout, null);
}

public TransactionalTopologyBuilder(String id, String spoutId, IPartitionedTransactionalSpout spout, Integer spoutParallelism) {
this(id, spoutId, new PartitionedTransactionalSpoutExecutor(spout), spoutParallelism);
}

public TransactionalTopologyBuilder(String id, String spoutId, ITransactionalSpout spout) {
public TransactionalTopologyBuilder(String id, String spoutId, IPartitionedTransactionalSpout spout) {
this(id, spoutId, spout, null);
}

public TransactionalTopologyBuilder(String id, String spoutId, IOpaquePartitionedTransactionalSpout spout, Integer spoutParallelism) {
this(id, spoutId, new OpaquePartitionedTransactionalSpoutExecutor(spout), spoutParallelism);
}

public TransactionalTopologyBuilder(String id, String spoutId, IPartitionedTransactionalSpout spout) {
public TransactionalTopologyBuilder(String id, String spoutId, IOpaquePartitionedTransactionalSpout spout) {
this(id, spoutId, spout, null);
}


public SpoutDeclarer getSpoutDeclarer() {
return new SpoutDeclarerImpl();
Expand Down
Expand Up @@ -84,7 +84,7 @@ public void commit(TransactionAttempt attempt) {
Object meta = metas.get(partition);
_partitionStates.get(partition).overrideState(txid, meta);
}
}
}

@Override
public void close() {
Expand Down

0 comments on commit 3df25ab

Please sign in to comment.