add kafka test conf

This commit is contained in:
Zhang Jianming 2022-06-23 22:46:13 +08:00
parent e5817f83d7
commit d20985a5cb
3 changed files with 13 additions and 13 deletions

View File

@ -108,7 +108,7 @@ public class KafkaConnectorConfig
/** /**
* whether to use kerberos * whether to use kerberos
*/ */
private boolean kerberosOn; private String kerberosOn;
public String getKrb5Conf() public String getKrb5Conf()
{ {
@ -206,7 +206,7 @@ public class KafkaConnectorConfig
return this; return this;
} }
public boolean isKerberosOn() public String isKerberosOn()
{ {
return kerberosOn; return kerberosOn;
} }
@ -216,7 +216,7 @@ public class KafkaConnectorConfig
defaultValue = "false", defaultValue = "false",
required = true) required = true)
@Config("kerberos.on") @Config("kerberos.on")
public KafkaConnectorConfig setKerberosOn(boolean kerberosOn) public KafkaConnectorConfig setKerberosOn(String kerberosOn)
{ {
this.kerberosOn = kerberosOn; this.kerberosOn = kerberosOn;
return this; return this;

View File

@ -47,7 +47,7 @@ public class KafkaSimpleConsumerManager
private final int connectTimeoutMillis; private final int connectTimeoutMillis;
private final int bufferSizeBytes; private final int bufferSizeBytes;
private final boolean kerberosOn; private final String kerberosOn;
private final String loginConfig; private final String loginConfig;
private final String krb5Conf; private final String krb5Conf;
private String groupId; private String groupId;
@ -116,7 +116,7 @@ public class KafkaSimpleConsumerManager
{ {
log.info("Creating new SaslConsumer for %s", host); log.info("Creating new SaslConsumer for %s", host);
Properties props = new Properties(); Properties props = new Properties();
if (kerberosOn) { if ("true".equalsIgnoreCase(kerberosOn)) {
System.setProperty("java.security.auth.login.config", loginConfig); System.setProperty("java.security.auth.login.config", loginConfig);
System.setProperty("java.security.krb5.conf", krb5Conf); System.setProperty("java.security.krb5.conf", krb5Conf);
props.put("security.protocol", securityProtocol); props.put("security.protocol", securityProtocol);

View File

@ -32,13 +32,13 @@ public class TestKafkaConnectorConfig
.setDefaultSchema("default") .setDefaultSchema("default")
.setTableNames("") .setTableNames("")
.setTableDescriptionDir(new File("etc/kafka/")) .setTableDescriptionDir(new File("etc/kafka/"))
.setGroupId("test") .setGroupId(null)
.setKerberosOn(false) .setKerberosOn(null)
.setSecurityProtocol("SASL_PLAINTEXT") .setSecurityProtocol(null)
.setKrb5Conf("/etc/krb5.conf") .setKrb5Conf(null)
.setLoginConfig("/etc/kafka_client_jaas.conf") .setLoginConfig(null)
.setSaslKerberosServiceName("kafka") .setSaslKerberosServiceName(null)
.setSaslMechanism("GSSAPI") .setSaslMechanism(null)
.setHideInternalColumns(true)); .setHideInternalColumns(true));
} }
@ -74,7 +74,7 @@ public class TestKafkaConnectorConfig
.setLoginConfig("/etc/kafka_client_jaas.conf") .setLoginConfig("/etc/kafka_client_jaas.conf")
.setSaslKerberosServiceName("kafka") .setSaslKerberosServiceName("kafka")
.setSaslMechanism("GSSAPI") .setSaslMechanism("GSSAPI")
.setKerberosOn(false) .setKerberosOn("false")
.setSecurityProtocol("SASL_PLAINTEXT") .setSecurityProtocol("SASL_PLAINTEXT")
.setHideInternalColumns(false); .setHideInternalColumns(false);