public class KafkaSource extends AbstractModuleFixture<KafkaSource>
| 修飾子と型 | フィールドと説明 |
|---|---|
static java.lang.String | DEFAULT_OUTPUT_TYPE |
static java.lang.String | DEFAULT_TOPIC |
static java.lang.String | DEFAULT_ZK_CLIENT |
label| コンストラクターと説明 |
|---|
KafkaSource(java.lang.String zkConnect)KafkaSource フィクスチャを初期化します。 |
| 修飾子と型 | メソッドと説明 |
|---|---|
KafkaSource | ensureReady()Zookeeper ソケットが最大 2 秒間ポーリングして使用可能であることを確認し、このソースに必要なトピックを作成します。 |
KafkaSource | outputType(java.lang.String outputType)kafka ソースの outputType を設定する |
protected java.lang.String | toDSL() ストリーム定義に含めるのに適したモジュールの表現を返します。 例: file --dir=xxxx --name=yyyy |
KafkaSource | topic(java.lang.String topic)kafka ソースのトピックを設定します |
static KafkaSource | withDefaults() デフォルトを使用して KafkaSource のインスタンスを返します。 |
KafkaSource | zkConnect(java.lang.String zkConnect)kafka ソースの zkConnect を設定する |
label, toStringpublic static final java.lang.String DEFAULT_ZK_CLIENT
public static final java.lang.String DEFAULT_TOPIC
public static final java.lang.String DEFAULT_OUTPUT_TYPE
public KafkaSource(java.lang.String zkConnect)
zkConnect - Zookeeper の接続文字列。public static KafkaSource withDefaults()
protected java.lang.String toDSL()
AbstractModuleFixturefile --dir=xxxx --name=yyyyAbstractModuleFixture<KafkaSource> の toDSL public KafkaSource topic(java.lang.String topic)
topic - データが投稿されるトピック。public KafkaSource zkConnect(java.lang.String zkConnect)
zkConnect - 使用する Zookeeper 接続文字列 public KafkaSource outputType(java.lang.String outputType)
outputType - 使用する出力型。public KafkaSource ensureReady()
java.lang.IllegalStateException - 2 秒以内に接続できない場合。