public class KafkaStarterUtils extends Object
| Modifier and Type | Field and Description |
|---|---|
static String |
BROKER_ID |
static int |
DEFAULT_BROKER_ID |
static String |
DEFAULT_KAFKA_BROKER |
static int |
DEFAULT_KAFKA_PORT |
static String |
DEFAULT_ZK_STR |
static int |
DEFAULT_ZK_TEST_PORT |
static String |
KAFKA_JSON_MESSAGE_DECODER_CLASS_NAME |
static String |
KAFKA_PRODUCER_CLASS_NAME |
static String |
KAFKA_SERVER_STARTABLE_CLASS_NAME |
static String |
KAFKA_STREAM_CONSUMER_FACTORY_CLASS_NAME |
static String |
KAFKA_STREAM_LEVEL_CONSUMER_CLASS_NAME |
static String |
PORT |
static String |
ZOOKEEPER_CONNECT |
| Constructor and Description |
|---|
KafkaStarterUtils() |
| Modifier and Type | Method and Description |
|---|---|
static void |
configureHostName(Properties configuration,
String hostName) |
static void |
configureOffsetsTopicReplicationFactor(Properties configuration,
short replicationFactor) |
static void |
configureTopicDeletion(Properties configuration,
boolean topicDeletionEnabled) |
static Properties |
getDefaultKafkaConfiguration() |
static Properties |
getTopicCreationProps(int numKafkaPartitions) |
static org.apache.pinot.spi.stream.StreamDataServerStartable |
startServer(int port,
int brokerId,
String zkStr,
Properties configuration) |
static List<org.apache.pinot.spi.stream.StreamDataServerStartable> |
startServers(int brokerCount,
int port,
String zkStr,
Properties configuration) |
public static final int DEFAULT_BROKER_ID
public static final int DEFAULT_ZK_TEST_PORT
public static final String DEFAULT_ZK_STR
public static int DEFAULT_KAFKA_PORT
public static final String DEFAULT_KAFKA_BROKER
public static final String PORT
public static final String BROKER_ID
public static final String ZOOKEEPER_CONNECT
public static final String KAFKA_SERVER_STARTABLE_CLASS_NAME
public static final String KAFKA_PRODUCER_CLASS_NAME
public static final String KAFKA_STREAM_CONSUMER_FACTORY_CLASS_NAME
public static final String KAFKA_STREAM_LEVEL_CONSUMER_CLASS_NAME
public static final String KAFKA_JSON_MESSAGE_DECODER_CLASS_NAME
public static Properties getDefaultKafkaConfiguration()
public static void configureOffsetsTopicReplicationFactor(Properties configuration, short replicationFactor)
public static void configureTopicDeletion(Properties configuration, boolean topicDeletionEnabled)
public static void configureHostName(Properties configuration, String hostName)
public static Properties getTopicCreationProps(int numKafkaPartitions)
public static List<org.apache.pinot.spi.stream.StreamDataServerStartable> startServers(int brokerCount, int port, String zkStr, Properties configuration)
public static org.apache.pinot.spi.stream.StreamDataServerStartable startServer(int port,
int brokerId,
String zkStr,
Properties configuration)
Copyright © 2018–2020 Apache Software Foundation. All rights reserved.