Class KafkaStarterUtils


  • public class KafkaStarterUtils
    extends Object
    • Field Detail

      • KAFKA_SERVER_STARTABLE_CLASS_NAME

        public static final String KAFKA_SERVER_STARTABLE_CLASS_NAME
      • KAFKA_PRODUCER_CLASS_NAME

        public static final String KAFKA_PRODUCER_CLASS_NAME
      • KAFKA_STREAM_CONSUMER_FACTORY_CLASS_NAME

        public static final String KAFKA_STREAM_CONSUMER_FACTORY_CLASS_NAME
      • KAFKA_STREAM_LEVEL_CONSUMER_CLASS_NAME

        public static final String KAFKA_STREAM_LEVEL_CONSUMER_CLASS_NAME
      • KAFKA_JSON_MESSAGE_DECODER_CLASS_NAME

        public static final String KAFKA_JSON_MESSAGE_DECODER_CLASS_NAME
        See Also:
        Constant Field Values
    • Method Detail

      • getDefaultKafkaConfiguration

        public static Properties getDefaultKafkaConfiguration()
      • configureOffsetsTopicReplicationFactor

        public static void configureOffsetsTopicReplicationFactor​(Properties configuration,
                                                                  short replicationFactor)
      • configureTransactionStateLogReplicationFactor

        public static void configureTransactionStateLogReplicationFactor​(Properties configuration,
                                                                         short replicationFactor)
      • configureTopicDeletion

        public static void configureTopicDeletion​(Properties configuration,
                                                  boolean topicDeletionEnabled)
      • configureHostName

        public static void configureHostName​(Properties configuration,
                                             String hostName)
      • getDefaultKafkaZKAddress

        public static String getDefaultKafkaZKAddress()
      • getTopicCreationProps

        public static Properties getTopicCreationProps​(int numKafkaPartitions)
      • startServers

        public static List<org.apache.pinot.spi.stream.StreamDataServerStartable> startServers​(int brokerCount,
                                                                                               int port,
                                                                                               String zkStr,
                                                                                               Properties configuration)
      • startServer

        public static org.apache.pinot.spi.stream.StreamDataServerStartable startServer​(int port,
                                                                                        int brokerId,
                                                                                        String zkStr,
                                                                                        Properties baseConf)