diff --git a/hetu-docs/zh/connector/kafka.md b/hetu-docs/zh/connector/kafka.md index ab1f61156..0adf0eaca 100644 --- a/hetu-docs/zh/connector/kafka.md +++ b/hetu-docs/zh/connector/kafka.md @@ -29,16 +29,22 @@ kafka.nodes=host1:port,host2:port 配置属性包括: -| 属性名称| 说明| -|:----------|:----------| -| `kafka.table-names`| 目录提供的所有表列表| -| `kafka.default-schema`| 表的默认模式名| -| `kafka.nodes`| Kafka集群节点列表| -| `kafka.connect-timeout`| 连接Kafka集群超时| -| `kafka.buffer-size`| Kafka读缓冲区大小| -| `kafka.table-description-dir`| 包含主题描述文件的目录| -| `kafka.hide-internal-columns`| 控制内部列是否是表模式的一部分| - +| 属性名称| 说明 | +|:----------|:-----------------------------------| +| `kafka.table-names`| 目录提供的所有表列表 | +| `kafka.default-schema`| 表的默认模式名 | +| `kafka.nodes`| Kafka集群节点列表 | +| `kafka.connect-timeout`| 连接Kafka集群超时 | +| `kafka.buffer-size`| Kafka读缓冲区大小 | +| `kafka.table-description-dir`| 包含主题描述文件的目录 | +| `kafka.hide-internal-columns`| 控制内部列是否是表模式的一部分 | +| `kerberos.on`| 是否开启Kerberos认证 | +| `java.security.auth.login.config`| kafka_client_jass.conf路径 | +| `java.security.krb5.conf`| krb5.conf文件路径 | +| `group.id`| kafka的groupID | +| `security.protocol`| Kafka的安全认证协议 | +| `sasl.mechanism`| sasl机制 | +| `sasl.kerberos.service.name`| kafka服务运行时的kerberos principal name | ### `kafka.table-names` 此目录提供的所有表的逗号分隔列表。表名可以是非限定的(简单名称),并将被放入默认模式(见下文)中,或者用模式名称(`.`)限定。 @@ -87,6 +93,48 @@ openLooKeng必须仍然能够连接到群集的所有节点,即使这里只指 此属性是可选的;默认值为`true`。 +### `kerberos.on` + +是否开启kerberos认证,适用于开启了kerberos认证的集群,如果在运行presto-kafka中的测试包,请置为false,因为测试程序使用内嵌Kafka,不支持认证。 + +此属性是可选的;默认值为`false`。 + +### `java.security.auth.login.config` + +Kafka的jaas_conf路径,也就是java认证和授权的相关文件,文件中存放的是认证和授权信息。 + +此属性是可选的;默认值为``。 + +### `java.security.krb5.conf` + +存放krb5.conf文件的路径,要注意全局配置中也需要配置此选项,例如部署后在jvm.config中配置,而在开发中需要在启动PrestoServer时使用"-D"参数配置。 + +此属性是可选的;默认值为``。 + +### `group.id` + +kafka的groupId。 + +此属性是可选的;默认值为``。 + +### `security.protocol` + +Kafka的安全协议。 + +此属性是可选的;默认值为`SASL_PLAINTEXT`。 + +### `sasl.mechanism` + +sasl机制,被用于客户端连接安全的机制。 + +此属性是可选的;默认值为`GSSAPI`。 + +### `sasl.kerberos.service.name` + +kafka运行时的kerberos principal name。 + +此属性是可选的;默认值为`kafka`。 + ## 内部列 对于每个已定义的表,连接器维护以下列: diff --git a/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaConnectorConfig.java b/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaConnectorConfig.java index 4a93d5e40..ae5708ac8 100644 --- a/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaConnectorConfig.java +++ b/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaConnectorConfig.java @@ -70,6 +70,158 @@ public class KafkaConnectorConfig */ private boolean hideInternalColumns = true; + /** + * the path of krb5.conf ,used for develop + */ + private String krb5Conf; + + /** + * the path of kafka_client_jaas_conf + */ + private String loginConfig; + + /** + * whether use subject creds only + */ + private String useSubjectCredsOnly; + + /** + * the group id of kafka + */ + private String groupId; + + /** + * the security protocol of kafka + */ + private String securityProtocol; + + /** + * the sasl mechanism of kafka + */ + private String saslMechanism; + + /** + * the sasl kerberos service name of kafka + */ + private String saslKerberosServiceName; + + /** + * whether to use kerberos + */ + private String kerberosOn; + + public String getKrb5Conf() + { + return krb5Conf; + } + + @Mandatory(name = "java.security.krb5.conf", + description = "java.security.krb5.conf", + defaultValue = "", + required = false) + @Config("java.security.krb5.conf") + public KafkaConnectorConfig setKrb5Conf(String krb5Conf) + { + this.krb5Conf = krb5Conf; + return this; + } + + public String getLoginConfig() + { + return loginConfig; + } + + @Mandatory(name = "java.security.auth.login.config", + description = "java.security.auth.login.config", + defaultValue = "", + required = false) + @Config("java.security.auth.login.config") + public KafkaConnectorConfig setLoginConfig(String loginConfig) + { + this.loginConfig = loginConfig; + return this; + } + + public String getGroupId() + { + return groupId; + } + + @Mandatory(name = "group.id", + description = "group.id", + defaultValue = "test", + required = false) + @Config("group.id") + public KafkaConnectorConfig setGroupId(String groupId) + { + this.groupId = groupId; + return this; + } + + public String getSecurityProtocol() + { + return securityProtocol; + } + + @Mandatory(name = "security.protocol", + description = "security.protocol", + defaultValue = "SASL_PLAINTEXT", + required = false) + @Config("security.protocol") + public KafkaConnectorConfig setSecurityProtocol(String securityProtocol) + { + this.securityProtocol = securityProtocol; + return this; + } + + public String getSaslMechanism() + { + return saslMechanism; + } + + @Mandatory(name = "sasl.mechanism", + description = "sasl.mechanism", + defaultValue = "GSSAPI", + required = false) + @Config("sasl.mechanism") + public KafkaConnectorConfig setSaslMechanism(String saslMechanism) + { + this.saslMechanism = saslMechanism; + return this; + } + + public String getSaslKerberosServiceName() + { + return saslKerberosServiceName; + } + + @Mandatory(name = "sasl.kerberos.service.name", + description = "sasl.kerberos.service.name", + defaultValue = "kafka", + required = false) + @Config("sasl.kerberos.service.name") + public KafkaConnectorConfig setSaslKerberosServiceName(String saslKerberosServiceName) + { + this.saslKerberosServiceName = saslKerberosServiceName; + return this; + } + + public String isKerberosOn() + { + return kerberosOn; + } + + @Mandatory(name = "kerberos.on", + description = "whether to use kerberos", + defaultValue = "false", + required = true) + @Config("kerberos.on") + public KafkaConnectorConfig setKerberosOn(String kerberosOn) + { + this.kerberosOn = kerberosOn; + return this; + } + @NotNull public File getTableDescriptionDir() { diff --git a/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaRecordSet.java b/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaRecordSet.java index 5b612ff71..30b1e86d1 100644 --- a/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaRecordSet.java +++ b/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaRecordSet.java @@ -27,11 +27,13 @@ import io.prestosql.spi.connector.RecordSet; import io.prestosql.spi.type.Type; import kafka.api.FetchRequest; import kafka.api.FetchRequestBuilder; -import kafka.javaapi.FetchResponse; -import kafka.javaapi.consumer.SimpleConsumer; -import kafka.message.MessageAndOffset; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.common.TopicPartition; import java.nio.ByteBuffer; +import java.util.Collections; import java.util.HashMap; import java.util.Iterator; import java.util.List; @@ -109,9 +111,9 @@ public class KafkaRecordSet private long totalBytes; private long totalMessages; private long cursorOffset = split.getStart(); - private Iterator messageAndOffsetIterator; + private Iterator> recordIterator; private final AtomicBoolean reported = new AtomicBoolean(); - + private KafkaConsumer leaderKafkaConsumer; private final FieldValueProvider[] currentRowValues = new FieldValueProvider[columnHandles.size()]; KafkaRecordCursor() @@ -147,19 +149,19 @@ public class KafkaRecordSet // Create a fetch request openFetchRequest(); - while (messageAndOffsetIterator.hasNext()) { - MessageAndOffset currentMessageAndOffset = messageAndOffsetIterator.next(); - long messageOffset = currentMessageAndOffset.offset(); + while (recordIterator.hasNext()) { + ConsumerRecord record = recordIterator.next(); + long messageOffset = record.offset(); if (messageOffset >= split.getEnd()) { return endOfData(); // Past our split end. Bail. } if (messageOffset >= cursorOffset) { - return nextRow(currentMessageAndOffset); + return nextRow(record); } } - messageAndOffsetIterator = null; + recordIterator = null; } } @@ -173,21 +175,21 @@ public class KafkaRecordSet return false; } - private boolean nextRow(MessageAndOffset messageAndOffset) + private boolean nextRow(ConsumerRecord record) { - cursorOffset = messageAndOffset.offset() + 1; // Cursor now points to the next message. - totalBytes += messageAndOffset.message().payloadSize(); + cursorOffset = record.offset() + 1; // Cursor now points to the next message. + totalBytes += record.serializedValueSize(); totalMessages++; byte[] keyData = EMPTY_BYTE_ARRAY; byte[] messageData = EMPTY_BYTE_ARRAY; - ByteBuffer key = messageAndOffset.message().key(); + ByteBuffer key = record.key(); if (key != null) { keyData = new byte[key.remaining()]; key.get(keyData); } - ByteBuffer message = messageAndOffset.message().payload(); + ByteBuffer message = record.value(); if (message != null) { messageData = new byte[message.remaining()]; message.get(messageData); @@ -206,7 +208,7 @@ public class KafkaRecordSet currentRowValuesMap.put(columnHandle, longValueProvider(totalMessages)); break; case PARTITION_OFFSET_FIELD: - currentRowValuesMap.put(columnHandle, longValueProvider(messageAndOffset.offset())); + currentRowValuesMap.put(columnHandle, longValueProvider(record.offset())); break; case MESSAGE_FIELD: currentRowValuesMap.put(columnHandle, bytesValueProvider(messageData)); @@ -305,12 +307,15 @@ public class KafkaRecordSet @Override public void close() { + if (leaderKafkaConsumer != null) { + leaderKafkaConsumer.close(); + } } private void openFetchRequest() { try { - if (messageAndOffsetIterator == null) { + if (recordIterator == null) { log.debug("Fetching %d bytes from offset %d (%d - %d). %d messages read so far", KAFKA_READ_BUFFER_SIZE, cursorOffset, split.getStart(), split.getEnd(), totalMessages); FetchRequest req = new FetchRequestBuilder() .clientId("presto-worker-" + Thread.currentThread().getName()) @@ -319,16 +324,14 @@ public class KafkaRecordSet // TODO - this should look at the actual node this is running on and prefer // that copy if running locally. - look into NodeInfo - SimpleConsumer consumer = consumerManager.getConsumer(split.getLeader()); - - FetchResponse fetchResponse = consumer.fetch(req); - if (fetchResponse.hasError()) { - short errorCode = fetchResponse.errorCode(split.getTopicName(), split.getPartitionId()); - log.warn("Fetch response has error: %d", errorCode); - throw new RuntimeException("could not fetch data from Kafka, error code is '" + errorCode + "'"); + if (leaderKafkaConsumer == null) { + leaderKafkaConsumer = consumerManager.getSaslConsumer(split.getLeader()); } - - messageAndOffsetIterator = fetchResponse.messageSet(split.getTopicName(), split.getPartitionId()).iterator(); + TopicPartition topicPartition = new TopicPartition(split.getTopicName(), split.getPartitionId()); + leaderKafkaConsumer.assign(Collections.singletonList(topicPartition)); + leaderKafkaConsumer.seek(topicPartition, cursorOffset); + ConsumerRecords records = leaderKafkaConsumer.poll(500); + recordIterator = records.records(topicPartition).iterator(); } } catch (Exception e) { // Catch all exceptions because Kafka library is written in scala and checked exceptions are not declared in method signature. diff --git a/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaSimpleConsumerManager.java b/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaSimpleConsumerManager.java index 1f082360b..91b5ccc3f 100644 --- a/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaSimpleConsumerManager.java +++ b/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaSimpleConsumerManager.java @@ -20,11 +20,15 @@ import io.airlift.log.Logger; import io.prestosql.spi.HostAddress; import io.prestosql.spi.NodeManager; import kafka.javaapi.consumer.SimpleConsumer; +import org.apache.kafka.clients.consumer.KafkaConsumer; import javax.annotation.PreDestroy; import javax.inject.Inject; +import java.nio.ByteBuffer; import java.util.Map; +import java.util.Properties; +import java.util.UUID; import static java.lang.Math.toIntExact; import static java.util.Objects.requireNonNull; @@ -43,6 +47,14 @@ public class KafkaSimpleConsumerManager private final int connectTimeoutMillis; private final int bufferSizeBytes; + private final String kerberosOn; + private final String loginConfig; + private final String krb5Conf; + private String groupId; + private final String securityProtocol; + private final String saslMechanism; + private final String saslKerberosServiceName; + @Inject public KafkaSimpleConsumerManager( KafkaConnectorConfig kafkaConnectorConfig, @@ -55,6 +67,14 @@ public class KafkaSimpleConsumerManager this.bufferSizeBytes = toIntExact(kafkaConnectorConfig.getKafkaBufferSize().toBytes()); this.consumerCache = CacheBuilder.newBuilder().build(CacheLoader.from(this::createConsumer)); + + this.kerberosOn = kafkaConnectorConfig.isKerberosOn(); + this.loginConfig = kafkaConnectorConfig.getLoginConfig(); + this.krb5Conf = kafkaConnectorConfig.getKrb5Conf(); + this.groupId = kafkaConnectorConfig.getGroupId(); + this.securityProtocol = kafkaConnectorConfig.getSecurityProtocol(); + this.saslMechanism = kafkaConnectorConfig.getSaslMechanism(); + this.saslKerberosServiceName = kafkaConnectorConfig.getSaslKerberosServiceName(); } @PreDestroy @@ -76,6 +96,12 @@ public class KafkaSimpleConsumerManager return consumerCache.getUnchecked(host); } + public KafkaConsumer getSaslConsumer(HostAddress host) + { + requireNonNull(host, "host is null"); + return createSaslConsumer(host); + } + private SimpleConsumer createConsumer(HostAddress host) { log.info("Creating new Consumer for %s", host); @@ -85,4 +111,35 @@ public class KafkaSimpleConsumerManager bufferSizeBytes, "presto-kafka-" + nodeManager.getCurrentNode().getNodeIdentifier()); } + + private KafkaConsumer createSaslConsumer(HostAddress host) + { + log.info("Creating new SaslConsumer for %s", host); + Properties props = new Properties(); + if ("true".equalsIgnoreCase(kerberosOn)) { + System.setProperty("java.security.auth.login.config", loginConfig); + System.setProperty("java.security.krb5.conf", krb5Conf); + props.put("security.protocol", securityProtocol); + props.put("sasl.mechanism", saslMechanism); + props.put("sasl.kerberos.service.name", saslKerberosServiceName); + } + + try { + props.put("bootstrap.servers", host.toString()); + props.put("enable.auto.commit", "false"); + props.put("key.deserializer", Class.forName("org.apache.kafka.common.serialization.ByteBufferDeserializer")); + props.put("value.deserializer", Class.forName("org.apache.kafka.common.serialization.ByteBufferDeserializer")); + if (groupId == null) { + groupId = UUID.randomUUID().toString(); + } + props.put("group.id", groupId); + props.put("session.timeout.ms", connectTimeoutMillis); + props.put("receive.buffer.bytes", bufferSizeBytes); + } + catch (ClassNotFoundException e) { + log.error(e, "failed to create kafka consumer"); + } + + return new KafkaConsumer<>(props); + } } diff --git a/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaSplitManager.java b/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaSplitManager.java index be74c05c3..eb480cbcc 100644 --- a/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaSplitManager.java +++ b/presto-kafka/src/main/java/io/prestosql/plugin/kafka/KafkaSplitManager.java @@ -28,15 +28,14 @@ import io.prestosql.spi.connector.ConnectorTableHandle; import io.prestosql.spi.connector.ConnectorTransactionHandle; import io.prestosql.spi.connector.FixedSplitSource; import kafka.api.PartitionOffsetRequestInfo; -import kafka.cluster.BrokerEndPoint; import kafka.common.TopicAndPartition; import kafka.javaapi.OffsetRequest; import kafka.javaapi.OffsetResponse; -import kafka.javaapi.PartitionMetadata; -import kafka.javaapi.TopicMetadata; -import kafka.javaapi.TopicMetadataRequest; -import kafka.javaapi.TopicMetadataResponse; import kafka.javaapi.consumer.SimpleConsumer; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.common.Node; +import org.apache.kafka.common.PartitionInfo; +import org.apache.kafka.common.TopicPartition; import javax.inject.Inject; @@ -47,6 +46,7 @@ import java.io.InputStreamReader; import java.net.MalformedURLException; import java.net.URI; import java.net.URL; +import java.nio.ByteBuffer; import java.util.List; import java.util.Set; import java.util.concurrent.ThreadLocalRandom; @@ -84,44 +84,30 @@ public class KafkaSplitManager public ConnectorSplitSource getSplits(ConnectorTransactionHandle transaction, ConnectorSession session, ConnectorTableHandle table, SplitSchedulingStrategy splitSchedulingStrategy) { KafkaTableHandle kafkaTableHandle = (KafkaTableHandle) table; - try { - SimpleConsumer simpleConsumer = consumerManager.getConsumer(selectRandom(nodes)); - - TopicMetadataRequest topicMetadataRequest = new TopicMetadataRequest(ImmutableList.of(kafkaTableHandle.getTopicName())); - TopicMetadataResponse topicMetadataResponse = simpleConsumer.send(topicMetadataRequest); + try (KafkaConsumer kafkaConsumer = consumerManager.getSaslConsumer(selectRandom(nodes))) { + List partitionInfos = kafkaConsumer.partitionsFor(kafkaTableHandle.getTopicName()); ImmutableList.Builder splits = ImmutableList.builder(); - for (TopicMetadata metadata : topicMetadataResponse.topicsMetadata()) { - for (PartitionMetadata part : metadata.partitionsMetadata()) { - log.debug("Adding Partition %s/%s", metadata.topic(), part.partitionId()); - - BrokerEndPoint leader = part.leader(); - if (leader == null) { - throw new PrestoException(GENERIC_INTERNAL_ERROR, format("Leader election in progress for Kafka topic '%s' partition %s", metadata.topic(), part.partitionId())); - } - - HostAddress partitionLeader = HostAddress.fromParts(leader.host(), leader.port()); - - SimpleConsumer leaderConsumer = consumerManager.getConsumer(partitionLeader); - // Kafka contains a reverse list of "end - start" pairs for the splits - - long[] offsets = findAllOffsets(leaderConsumer, metadata.topic(), part.partitionId()); - - for (int i = offsets.length - 1; i > 0; i--) { - KafkaSplit split = new KafkaSplit( - metadata.topic(), - kafkaTableHandle.getKeyDataFormat(), - kafkaTableHandle.getMessageDataFormat(), - kafkaTableHandle.getKeyDataSchemaLocation().map(KafkaSplitManager::readSchema), - kafkaTableHandle.getMessageDataSchemaLocation().map(KafkaSplitManager::readSchema), - part.partitionId(), - offsets[i], - offsets[i - 1], - partitionLeader); - splits.add(split); - } - } + for (PartitionInfo partitionInfo : partitionInfos) { + log.debug("Adding Partition %s/%s", partitionInfo.topic(), partitionInfo.partition()); + Node leader = partitionInfo.leader(); + HostAddress partitionLeader = HostAddress.fromParts(leader.host(), leader.port()); + TopicPartition topicPartition = new TopicPartition(partitionInfo.topic(), partitionInfo.partition()); + kafkaConsumer.assign(ImmutableList.of(topicPartition)); + long beginOffset = kafkaConsumer.beginningOffsets(ImmutableList.of(topicPartition)).values().iterator().next(); + long endOffset = kafkaConsumer.endOffsets(ImmutableList.of(topicPartition)).values().iterator().next(); + KafkaSplit split = new KafkaSplit( + topicPartition.topic(), + kafkaTableHandle.getKeyDataFormat(), + kafkaTableHandle.getMessageDataFormat(), + kafkaTableHandle.getKeyDataSchemaLocation().map(KafkaSplitManager::readSchema), + kafkaTableHandle.getMessageDataSchemaLocation().map(KafkaSplitManager::readSchema), + topicPartition.partition(), + beginOffset, + endOffset, + partitionLeader); + splits.add(split); } return new FixedSplitSource(splits.build()); diff --git a/presto-kafka/src/test/java/io/prestosql/plugin/kafka/TestKafkaConnectorConfig.java b/presto-kafka/src/test/java/io/prestosql/plugin/kafka/TestKafkaConnectorConfig.java index 27650796b..ac2d7444a 100644 --- a/presto-kafka/src/test/java/io/prestosql/plugin/kafka/TestKafkaConnectorConfig.java +++ b/presto-kafka/src/test/java/io/prestosql/plugin/kafka/TestKafkaConnectorConfig.java @@ -32,6 +32,13 @@ public class TestKafkaConnectorConfig .setDefaultSchema("default") .setTableNames("") .setTableDescriptionDir(new File("etc/kafka/")) + .setGroupId(null) + .setKerberosOn(null) + .setSecurityProtocol(null) + .setKrb5Conf(null) + .setLoginConfig(null) + .setSaslKerberosServiceName(null) + .setSaslMechanism(null) .setHideInternalColumns(true)); } @@ -46,6 +53,13 @@ public class TestKafkaConnectorConfig .put("kafka.connect-timeout", "1h") .put("kafka.buffer-size", "1MB") .put("kafka.hide-internal-columns", "false") + .put("group.id", "test") + .put("java.security.auth.login.config", "/etc/kafka_client_jaas.conf") + .put("java.security.krb5.conf", "/etc/krb5.conf") + .put("kerberos.on", "false") + .put("sasl.kerberos.service.name", "kafka") + .put("sasl.mechanism", "GSSAPI") + .put("security.protocol", "SASL_PLAINTEXT") .build(); KafkaConnectorConfig expected = new KafkaConnectorConfig() @@ -55,6 +69,13 @@ public class TestKafkaConnectorConfig .setNodes("localhost:12345, localhost:23456") .setKafkaConnectTimeout("1h") .setKafkaBufferSize("1MB") + .setGroupId("test") + .setKrb5Conf("/etc/krb5.conf") + .setLoginConfig("/etc/kafka_client_jaas.conf") + .setSaslKerberosServiceName("kafka") + .setSaslMechanism("GSSAPI") + .setKerberosOn("false") + .setSecurityProtocol("SASL_PLAINTEXT") .setHideInternalColumns(false); ConfigAssertions.assertFullMapping(properties, expected); diff --git a/presto-main/etc/catalog/dc.properties b/presto-main/etc/catalog/dc.properties deleted file mode 100644 index 51e333ffa..000000000 --- a/presto-main/etc/catalog/dc.properties +++ /dev/null @@ -1,4 +0,0 @@ -connector.name=dc -connection-url=http://localhost:8090 -connection-user=root -connection-password= \ No newline at end of file diff --git a/presto-main/etc/catalog/hive.properties b/presto-main/etc/catalog/hive.properties deleted file mode 100644 index 5f7221544..000000000 --- a/presto-main/etc/catalog/hive.properties +++ /dev/null @@ -1,9 +0,0 @@ -# -# WARNING -# ^^^^^^^ -# This configuration file is for development only and should NOT be used -# in production. For example configuration, see the Presto documentation. -# - -connector.name=hive-hadoop2 -hive.metastore.uri=thrift://localhost:9083 \ No newline at end of file diff --git a/presto-main/etc/catalog/kafka.properties b/presto-main/etc/catalog/kafka.properties new file mode 100644 index 000000000..c5025f6e6 --- /dev/null +++ b/presto-main/etc/catalog/kafka.properties @@ -0,0 +1,11 @@ +connector.name=kafka +kafka.nodes=localhost:9092 +kafka.table-names=testTopic +kafka.hide-internal-columns=false +kerberos.on=false +java.security.auth.login.config=/Users/path/kafka-jaas.conf +java.security.krb5.conf=/Users/path/krb5.conf +group.id=testTopic +security.protocol=SASL_PLAINTEXT +sasl.mechanism=GSSAPI +sasl.kerberos.service.name=kafka