Pre Merge pull request !1524 from 张剑鸣/kafka

This commit is contained in:
张剑鸣 2022-06-21 01:28:16 +00:00 committed by Gitee
commit 5d7e009882
No known key found for this signature in database
GPG Key ID: 173E9B9CA92EEF8F
6 changed files with 305 additions and 66 deletions

View File

@ -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`。
## 内部列
对于每个已定义的表,连接器维护以下列:

View File

@ -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()
{

View File

@ -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<MessageAndOffset> messageAndOffsetIterator;
private Iterator<ConsumerRecord<ByteBuffer, ByteBuffer>> recordIterator;
private final AtomicBoolean reported = new AtomicBoolean();
private KafkaConsumer<ByteBuffer, ByteBuffer> 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<ByteBuffer, ByteBuffer> 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<ByteBuffer, ByteBuffer> 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<ByteBuffer, ByteBuffer> 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.

View File

@ -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<ByteBuffer, ByteBuffer> 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<ByteBuffer, ByteBuffer> 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);
}
}

View File

@ -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<ByteBuffer, ByteBuffer> kafkaConsumer = consumerManager.getSaslConsumer(selectRandom(nodes))) {
List<PartitionInfo> partitionInfos = kafkaConsumer.partitionsFor(kafkaTableHandle.getTopicName());
ImmutableList.Builder<ConnectorSplit> 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());

View File

@ -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