diff --git a/hetu-docs/zh/connector/kafka.md b/hetu-docs/zh/connector/kafka.md index ab1f61156..dce59998d 100644 --- a/hetu-docs/zh/connector/kafka.md +++ b/hetu-docs/zh/connector/kafka.md @@ -87,6 +87,48 @@ openLooKeng必须仍然能够连接到群集的所有节点,即使这里只指 此属性是可选的;默认值为`true`。 +### `kerberos.on` + +是否开启kerberos认证,适用于开启了kerberos认证的集群。 + +此属性是可选的;默认值为`false`。 + +### `java.security.auth.login.config` + +Kafka的jaas_conf路径,也就是java认证和授权的相关文件,文件中存放的是认证和授权信息。 + +此属性是可选的;默认值为``。 + +### `java.security.krb5.conf` + +存放krb5.conf文件的路径,要注意全局配置中也需要配置此选项,例如部署后在jvm.config中配置。 + +此属性是可选的;默认值为``。 + +### `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..ae5cbee93 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,151 @@ 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 boolean 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 void setKrb5Conf(String krb5Conf) + { + this.krb5Conf = krb5Conf; + } + + 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 void setLoginConfig(String loginConfig) + { + this.loginConfig = loginConfig; + } + + public String getGroupId() + { + return groupId; + } + + @Mandatory(name = "group.id", + description = "group.id", + defaultValue = "", + required = false) + @Config("group.id") + public void setGroupId(String groupId) + { + this.groupId = groupId; + } + + public String getSecurityProtocol() + { + return securityProtocol; + } + + @Mandatory(name = "security.protocol", + description = "security.protocol", + defaultValue = "SASL_PLAINTEXT", + required = false) + @Config("security.protocol") + public void setSecurityProtocol(String securityProtocol) + { + this.securityProtocol = securityProtocol; + } + + public String getSaslMechanism() + { + return saslMechanism; + } + + @Mandatory(name = "sasl.mechanism", + description = "sasl.mechanism", + defaultValue = "GSSAPI", + required = false) + @Config("sasl.mechanism") + public void setSaslMechanism(String saslMechanism) + { + this.saslMechanism = saslMechanism; + } + + 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 void setSaslKerberosServiceName(String saslKerberosServiceName) + { + this.saslKerberosServiceName = saslKerberosServiceName; + } + + public boolean isKerberosOn() + { + return kerberosOn; + } + + @Mandatory(name = "kerberos.on", + description = "whether to use kerberos", + defaultValue = "false", + required = true) + @Config("kerberos.on") + public void setKerberosOn(boolean kerberosOn) + { + this.kerberosOn = kerberosOn; + } + @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..f3c70c04c 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,14 @@ 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 static java.lang.Math.toIntExact; import static java.util.Objects.requireNonNull; @@ -43,6 +46,14 @@ public class KafkaSimpleConsumerManager private final int connectTimeoutMillis; private final int bufferSizeBytes; + private final boolean kerberosOn; + private final String loginConfig; + private final String krb5Conf; + private final String groupId; + private final String securityProtocol; + private final String saslMechanism; + private final String saslKerberosServiceName; + @Inject public KafkaSimpleConsumerManager( KafkaConnectorConfig kafkaConnectorConfig, @@ -55,6 +66,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 +95,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 +110,30 @@ 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 (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")); + props.put("group.id", groupId); + } + catch (ClassNotFoundException e) { + e.printStackTrace(); + } + + 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-main/etc/catalog/kafka.properties b/presto-main/etc/catalog/kafka.properties new file mode 100644 index 000000000..eed0c8015 --- /dev/null +++ b/presto-main/etc/catalog/kafka.properties @@ -0,0 +1,12 @@ +connector.name=kafka +kafka.nodes=localhost:9092 +kafka.table-names=testTopic +kafka.hide-internal-columns=false +kerberos.on=true + +java.security.auth.login.config=/Users/mac/Desktop/kafka-jaas.conf +java.security.krb5.conf=/Users/mac/Desktop/krb5.conf +group.id=test1 +security.protocol=SASL_PLAINTEXT +sasl.mechanism=GSSAPI +sasl.kerberos.service.name=kafka