From 603f9e1a807bc2e084c79c28597cef260bb9b1eb Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 16 Jun 2022 21:54:29 +0800 Subject: [PATCH 01/12] update presto-Kafka for kerberos --- hetu-docs/zh/connector/kafka.md | 42 +++++ .../plugin/kafka/KafkaConnectorConfig.java | 145 ++++++++++++++++++ .../plugin/kafka/KafkaRecordSet.java | 55 +++---- .../kafka/KafkaSimpleConsumerManager.java | 51 ++++++ .../plugin/kafka/KafkaSplitManager.java | 66 ++++---- presto-main/etc/catalog/kafka.properties | 12 ++ 6 files changed, 305 insertions(+), 66 deletions(-) create mode 100644 presto-main/etc/catalog/kafka.properties 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 From 7798d50fabbb8bfea86d281d0338d066ab34e351 Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Tue, 21 Jun 2022 19:05:52 +0800 Subject: [PATCH 02/12] update presto-Kafka for kerberos --- hetu-docs/zh/connector/kafka.md | 26 ++++++++++++------- .../kafka/KafkaSimpleConsumerManager.java | 2 ++ 2 files changed, 18 insertions(+), 10 deletions(-) diff --git a/hetu-docs/zh/connector/kafka.md b/hetu-docs/zh/connector/kafka.md index dce59998d..985c6553a 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` 此目录提供的所有表的逗号分隔列表。表名可以是非限定的(简单名称),并将被放入默认模式(见下文)中,或者用模式名称(`.`)限定。 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 f3c70c04c..627fd2a04 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 @@ -129,6 +129,8 @@ public class KafkaSimpleConsumerManager 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); + props.put("session.timeout.ms", connectTimeoutMillis); + props.put("receive.buffer.bytes", bufferSizeBytes); } catch (ClassNotFoundException e) { e.printStackTrace(); From 3aff9171859a1cfca3ed092f4869108e3ad1ed90 Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 16:47:11 +0800 Subject: [PATCH 03/12] fix the null of group id --- hetu-docs/zh/connector/kafka.md | 2 +- .../plugin/kafka/KafkaConnectorConfig.java | 23 ++++++++++++------- .../kafka/KafkaSimpleConsumerManager.java | 8 +++++-- 3 files changed, 22 insertions(+), 11 deletions(-) diff --git a/hetu-docs/zh/connector/kafka.md b/hetu-docs/zh/connector/kafka.md index 985c6553a..b59d087f1 100644 --- a/hetu-docs/zh/connector/kafka.md +++ b/hetu-docs/zh/connector/kafka.md @@ -107,7 +107,7 @@ Kafka的jaas_conf路径,也就是java认证和授权的相关文件,文件 ### `java.security.krb5.conf` -存放krb5.conf文件的路径,要注意全局配置中也需要配置此选项,例如部署后在jvm.config中配置。 +存放krb5.conf文件的路径,要注意全局配置中也需要配置此选项,例如部署后在jvm.config中配置,而在开发中需要在启动PrestoServer时使用"-D"参数配置。 此属性是可选的;默认值为``。 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 ae5cbee93..1df2891ef 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 @@ -120,9 +120,10 @@ public class KafkaConnectorConfig defaultValue = "", required = false) @Config("java.security.krb5.conf") - public void setKrb5Conf(String krb5Conf) + public KafkaConnectorConfig setKrb5Conf(String krb5Conf) { this.krb5Conf = krb5Conf; + return this; } public String getLoginConfig() @@ -135,9 +136,10 @@ public class KafkaConnectorConfig defaultValue = "", required = false) @Config("java.security.auth.login.config") - public void setLoginConfig(String loginConfig) + public KafkaConnectorConfig setLoginConfig(String loginConfig) { this.loginConfig = loginConfig; + return this; } public String getGroupId() @@ -147,12 +149,13 @@ public class KafkaConnectorConfig @Mandatory(name = "group.id", description = "group.id", - defaultValue = "", + defaultValue = "test", required = false) @Config("group.id") - public void setGroupId(String groupId) + public KafkaConnectorConfig setGroupId(String groupId) { this.groupId = groupId; + return this; } public String getSecurityProtocol() @@ -165,9 +168,10 @@ public class KafkaConnectorConfig defaultValue = "SASL_PLAINTEXT", required = false) @Config("security.protocol") - public void setSecurityProtocol(String securityProtocol) + public KafkaConnectorConfig setSecurityProtocol(String securityProtocol) { this.securityProtocol = securityProtocol; + return this; } public String getSaslMechanism() @@ -180,9 +184,10 @@ public class KafkaConnectorConfig defaultValue = "GSSAPI", required = false) @Config("sasl.mechanism") - public void setSaslMechanism(String saslMechanism) + public KafkaConnectorConfig setSaslMechanism(String saslMechanism) { this.saslMechanism = saslMechanism; + return this; } public String getSaslKerberosServiceName() @@ -195,9 +200,10 @@ public class KafkaConnectorConfig defaultValue = "kafka", required = false) @Config("sasl.kerberos.service.name") - public void setSaslKerberosServiceName(String saslKerberosServiceName) + public KafkaConnectorConfig setSaslKerberosServiceName(String saslKerberosServiceName) { this.saslKerberosServiceName = saslKerberosServiceName; + return this; } public boolean isKerberosOn() @@ -210,9 +216,10 @@ public class KafkaConnectorConfig defaultValue = "false", required = true) @Config("kerberos.on") - public void setKerberosOn(boolean kerberosOn) + public KafkaConnectorConfig setKerberosOn(boolean kerberosOn) { this.kerberosOn = kerberosOn; + return this; } @NotNull 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 627fd2a04..056316072 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 @@ -28,6 +28,7 @@ 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; @@ -49,7 +50,7 @@ public class KafkaSimpleConsumerManager private final boolean kerberosOn; private final String loginConfig; private final String krb5Conf; - private final String groupId; + private String groupId; private final String securityProtocol; private final String saslMechanism; private final String saslKerberosServiceName; @@ -128,12 +129,15 @@ public class KafkaSimpleConsumerManager 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) { - e.printStackTrace(); + log.error(e, "failed to create kafka consumer"); } return new KafkaConsumer<>(props); From 1e9ed7ae29996759e4ba4ac8cb4515db108ab993 Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 16:53:44 +0800 Subject: [PATCH 04/12] fix the null of group id --- hetu-docs/zh/connector/kafka.md | 2 +- presto-main/etc/catalog/kafka.properties | 9 ++++----- 2 files changed, 5 insertions(+), 6 deletions(-) diff --git a/hetu-docs/zh/connector/kafka.md b/hetu-docs/zh/connector/kafka.md index b59d087f1..0adf0eaca 100644 --- a/hetu-docs/zh/connector/kafka.md +++ b/hetu-docs/zh/connector/kafka.md @@ -95,7 +95,7 @@ openLooKeng必须仍然能够连接到群集的所有节点,即使这里只指 ### `kerberos.on` -是否开启kerberos认证,适用于开启了kerberos认证的集群。 +是否开启kerberos认证,适用于开启了kerberos认证的集群,如果在运行presto-kafka中的测试包,请置为false,因为测试程序使用内嵌Kafka,不支持认证。 此属性是可选的;默认值为`false`。 diff --git a/presto-main/etc/catalog/kafka.properties b/presto-main/etc/catalog/kafka.properties index eed0c8015..c5025f6e6 100644 --- a/presto-main/etc/catalog/kafka.properties +++ b/presto-main/etc/catalog/kafka.properties @@ -2,11 +2,10 @@ 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 +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 From c2f6e746d9fcafcc1cec3e19c5712aadb328a3c0 Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 17:56:27 +0800 Subject: [PATCH 05/12] fix the null of group id --- .../prestosql/plugin/kafka/TestKafkaConnectorConfig.java | 4 +++- presto-main/etc/catalog/dc.properties | 4 ---- presto-main/etc/catalog/hive.properties | 9 --------- 3 files changed, 3 insertions(+), 14 deletions(-) delete mode 100644 presto-main/etc/catalog/dc.properties delete mode 100644 presto-main/etc/catalog/hive.properties 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..616f51c3d 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 @@ -31,7 +31,7 @@ public class TestKafkaConnectorConfig .setKafkaBufferSize("64kB") .setDefaultSchema("default") .setTableNames("") - .setTableDescriptionDir(new File("etc/kafka/")) + .setTableDescriptionDir(new File("etc/kafka/")).setGroupId("ccc") .setHideInternalColumns(true)); } @@ -46,6 +46,7 @@ public class TestKafkaConnectorConfig .put("kafka.connect-timeout", "1h") .put("kafka.buffer-size", "1MB") .put("kafka.hide-internal-columns", "false") + .put("group.id", "bbb") .build(); KafkaConnectorConfig expected = new KafkaConnectorConfig() @@ -55,6 +56,7 @@ public class TestKafkaConnectorConfig .setNodes("localhost:12345, localhost:23456") .setKafkaConnectTimeout("1h") .setKafkaBufferSize("1MB") + .setGroupId("aaa") .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 From 6a813fac58833143e2daff6172b569a088aef1e5 Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 18:23:37 +0800 Subject: [PATCH 06/12] add kafka test conf --- .../prestosql/plugin/kafka/TestKafkaConnectorConfig.java | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) 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 616f51c3d..8ca69a48b 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 @@ -31,7 +31,7 @@ public class TestKafkaConnectorConfig .setKafkaBufferSize("64kB") .setDefaultSchema("default") .setTableNames("") - .setTableDescriptionDir(new File("etc/kafka/")).setGroupId("ccc") + .setTableDescriptionDir(new File("etc/kafka/")) .setHideInternalColumns(true)); } @@ -46,7 +46,6 @@ public class TestKafkaConnectorConfig .put("kafka.connect-timeout", "1h") .put("kafka.buffer-size", "1MB") .put("kafka.hide-internal-columns", "false") - .put("group.id", "bbb") .build(); KafkaConnectorConfig expected = new KafkaConnectorConfig() @@ -56,7 +55,11 @@ public class TestKafkaConnectorConfig .setNodes("localhost:12345, localhost:23456") .setKafkaConnectTimeout("1h") .setKafkaBufferSize("1MB") - .setGroupId("aaa") + .setGroupId("test") + .setKrb5Conf("/etc/krb5.conf") + .setLoginConfig("/etc/kafka_client_jaas.conf") + .setSaslKerberosServiceName("kafka") + .setSaslMechanism("GSSAPI") .setHideInternalColumns(false); ConfigAssertions.assertFullMapping(properties, expected); From 2a9ed68d2cee0dd3ef6ce9370e6380b5385ec515 Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 18:39:07 +0800 Subject: [PATCH 07/12] add kafka test conf --- .../io/prestosql/plugin/kafka/TestKafkaConnectorConfig.java | 5 +++++ 1 file changed, 5 insertions(+) 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 8ca69a48b..c088f157d 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,11 @@ public class TestKafkaConnectorConfig .setDefaultSchema("default") .setTableNames("") .setTableDescriptionDir(new File("etc/kafka/")) + .setGroupId("test") + .setKrb5Conf("/etc/krb5.conf") + .setLoginConfig("/etc/kafka_client_jaas.conf") + .setSaslKerberosServiceName("kafka") + .setSaslMechanism("GSSAPI") .setHideInternalColumns(true)); } From f2664c261f212fd6bb887c2d3f03a17162f46fcc Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 18:56:50 +0800 Subject: [PATCH 08/12] add kafka test conf --- .../java/io/prestosql/plugin/kafka/TestKafkaConnectorConfig.java | 1 + 1 file changed, 1 insertion(+) 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 c088f157d..745f49926 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 @@ -51,6 +51,7 @@ public class TestKafkaConnectorConfig .put("kafka.connect-timeout", "1h") .put("kafka.buffer-size", "1MB") .put("kafka.hide-internal-columns", "false") + .put("group.id", "bbb") .build(); KafkaConnectorConfig expected = new KafkaConnectorConfig() From 7da529ee2b633806cf0969215b84ffc870d12814 Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 20:10:02 +0800 Subject: [PATCH 09/12] add kafka test conf --- .../plugin/kafka/TestKafkaConnectorConfig.java | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) 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 745f49926..d8b272457 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 @@ -51,7 +51,13 @@ public class TestKafkaConnectorConfig .put("kafka.connect-timeout", "1h") .put("kafka.buffer-size", "1MB") .put("kafka.hide-internal-columns", "false") - .put("group.id", "bbb") + .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() @@ -66,6 +72,8 @@ public class TestKafkaConnectorConfig .setLoginConfig("/etc/kafka_client_jaas.conf") .setSaslKerberosServiceName("kafka") .setSaslMechanism("GSSAPI") + .setKerberosOn(false) + .setSecurityProtocol("SASL_PLAINTEXT") .setHideInternalColumns(false); ConfigAssertions.assertFullMapping(properties, expected); From ca4d024b248276ee360e2ba2f326899eb92cb099 Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 20:30:27 +0800 Subject: [PATCH 10/12] add kafka test conf --- .../plugin/kafka/TestKafkaConnectorConfig.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) 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 d8b272457..7ecc62fdf 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 @@ -52,12 +52,12 @@ public class TestKafkaConnectorConfig .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") + .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() From e5817f83d79b949fddbcc5cb4250a234c16590ee Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 21:02:32 +0800 Subject: [PATCH 11/12] add kafka test conf --- .../io/prestosql/plugin/kafka/TestKafkaConnectorConfig.java | 2 ++ 1 file changed, 2 insertions(+) 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 7ecc62fdf..3300fa02a 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 @@ -33,6 +33,8 @@ public class TestKafkaConnectorConfig .setTableNames("") .setTableDescriptionDir(new File("etc/kafka/")) .setGroupId("test") + .setKerberosOn(false) + .setSecurityProtocol("SASL_PLAINTEXT") .setKrb5Conf("/etc/krb5.conf") .setLoginConfig("/etc/kafka_client_jaas.conf") .setSaslKerberosServiceName("kafka") From d20985a5cb429fac8aef78a1db80f5347c483932 Mon Sep 17 00:00:00 2001 From: Zhang Jianming Date: Thu, 23 Jun 2022 22:46:13 +0800 Subject: [PATCH 12/12] add kafka test conf --- .../plugin/kafka/KafkaConnectorConfig.java | 6 +++--- .../plugin/kafka/KafkaSimpleConsumerManager.java | 4 ++-- .../plugin/kafka/TestKafkaConnectorConfig.java | 16 ++++++++-------- 3 files changed, 13 insertions(+), 13 deletions(-) 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 1df2891ef..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 @@ -108,7 +108,7 @@ public class KafkaConnectorConfig /** * whether to use kerberos */ - private boolean kerberosOn; + private String kerberosOn; public String getKrb5Conf() { @@ -206,7 +206,7 @@ public class KafkaConnectorConfig return this; } - public boolean isKerberosOn() + public String isKerberosOn() { return kerberosOn; } @@ -216,7 +216,7 @@ public class KafkaConnectorConfig defaultValue = "false", required = true) @Config("kerberos.on") - public KafkaConnectorConfig setKerberosOn(boolean kerberosOn) + public KafkaConnectorConfig setKerberosOn(String kerberosOn) { this.kerberosOn = kerberosOn; return this; 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 056316072..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 @@ -47,7 +47,7 @@ public class KafkaSimpleConsumerManager private final int connectTimeoutMillis; private final int bufferSizeBytes; - private final boolean kerberosOn; + private final String kerberosOn; private final String loginConfig; private final String krb5Conf; private String groupId; @@ -116,7 +116,7 @@ public class KafkaSimpleConsumerManager { log.info("Creating new SaslConsumer for %s", host); Properties props = new Properties(); - if (kerberosOn) { + if ("true".equalsIgnoreCase(kerberosOn)) { System.setProperty("java.security.auth.login.config", loginConfig); System.setProperty("java.security.krb5.conf", krb5Conf); props.put("security.protocol", securityProtocol); 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 3300fa02a..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,13 +32,13 @@ public class TestKafkaConnectorConfig .setDefaultSchema("default") .setTableNames("") .setTableDescriptionDir(new File("etc/kafka/")) - .setGroupId("test") - .setKerberosOn(false) - .setSecurityProtocol("SASL_PLAINTEXT") - .setKrb5Conf("/etc/krb5.conf") - .setLoginConfig("/etc/kafka_client_jaas.conf") - .setSaslKerberosServiceName("kafka") - .setSaslMechanism("GSSAPI") + .setGroupId(null) + .setKerberosOn(null) + .setSecurityProtocol(null) + .setKrb5Conf(null) + .setLoginConfig(null) + .setSaslKerberosServiceName(null) + .setSaslMechanism(null) .setHideInternalColumns(true)); } @@ -74,7 +74,7 @@ public class TestKafkaConnectorConfig .setLoginConfig("/etc/kafka_client_jaas.conf") .setSaslKerberosServiceName("kafka") .setSaslMechanism("GSSAPI") - .setKerberosOn(false) + .setKerberosOn("false") .setSecurityProtocol("SASL_PLAINTEXT") .setHideInternalColumns(false);