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);