@SuppressWarnings({ "rawtypes", "unchecked" }) @SPI("kafka") public class CanalKafkaProducer extends AbstractMQProducer implements CanalMQProducer {
private static final Logger logger = LoggerFactory.getLogger(CanalKafkaProducer.class);
private static final String PREFIX_KAFKA_CONFIG = "kafka.";
private Producer<String, byte[]> producer;
@Override public void init(Properties properties) { KafkaProducerConfig kafkaProducerConfig = new KafkaProducerConfig(); this.mqProperties = kafkaProducerConfig; super.init(properties); this.loadKafkaProperties(properties);
Properties kafkaProperties = new Properties(); kafkaProperties.putAll(kafkaProducerConfig.getKafkaProperties()); kafkaProperties.put("max.in.flight.requests.per.connection", 1); kafkaProperties.put("key.serializer", StringSerializer.class); if (kafkaProducerConfig.isKerberosEnabled()) { File krb5File = new File(kafkaProducerConfig.getKrb5File()); File jaasFile = new File(kafkaProducerConfig.getJaasFile()); if (krb5File.exists() && jaasFile.exists()) { System.setProperty("java.security.krb5.conf", krb5File.getAbsolutePath()); System.setProperty("java.security.auth.login.config", jaasFile.getAbsolutePath()); System.setProperty("javax.security.auth.useSubjectCredsOnly", "false"); kafkaProperties.put("security.protocol", "SASL_PLAINTEXT"); kafkaProperties.put("sasl.kerberos.service.name", "kafka"); } else { String errorMsg = "ERROR # The kafka kerberos configuration file does not exist! please check it"; logger.error(errorMsg); throw new RuntimeException(errorMsg); } } kafkaProperties.put("value.serializer", KafkaMessageSerializer.class); producer = new KafkaProducer<>(kafkaProperties); }
private void loadKafkaProperties(Properties properties) { KafkaProducerConfig kafkaProducerConfig = (KafkaProducerConfig) this.mqProperties; Map<String, Object> kafkaProperties = kafkaProducerConfig.getKafkaProperties(); doMoreCompatibleConvert("canal.mq.servers", "kafka.bootstrap.servers", properties); doMoreCompatibleConvert("canal.mq.acks", "kafka.acks", properties); doMoreCompatibleConvert("canal.mq.compressionType", "kafka.compression.type", properties); doMoreCompatibleConvert("canal.mq.retries", "kafka.retries", properties); doMoreCompatibleConvert("canal.mq.batchSize", "kafka.batch.size", properties); doMoreCompatibleConvert("canal.mq.lingerMs", "kafka.linger.ms", properties); doMoreCompatibleConvert("canal.mq.maxRequestSize", "kafka.max.request.size", properties); doMoreCompatibleConvert("canal.mq.bufferMemory", "kafka.buffer.memory", properties); doMoreCompatibleConvert("canal.mq.kafka.kerberos.enable", "kafka.kerberos.enable", properties); doMoreCompatibleConvert("canal.mq.kafka.kerberos.krb5.file", "kafka.kerberos.krb5.file", properties); doMoreCompatibleConvert("canal.mq.kafka.kerberos.jaas.file", "kafka.kerberos.jaas.file", properties);
for (Map.Entry<Object, Object> entry : properties.entrySet()) { String key = (String) entry.getKey(); Object value = entry.getValue(); if (key.startsWith(PREFIX_KAFKA_CONFIG) && value != null) { value = PropertiesUtils.getProperty(properties, key); key = key.substring(PREFIX_KAFKA_CONFIG.length()); kafkaProperties.put(key, value); } } String kerberosEnabled = PropertiesUtils.getProperty(properties, KafkaConstants.CANAL_MQ_KAFKA_KERBEROS_ENABLE); if (!StringUtils.isEmpty(kerberosEnabled)) { kafkaProducerConfig.setKerberosEnabled(Boolean.parseBoolean(kerberosEnabled)); } String krb5File = PropertiesUtils.getProperty(properties, KafkaConstants.CANAL_MQ_KAFKA_KERBEROS_KRB5_FILE); if (!StringUtils.isEmpty(krb5File)) { kafkaProducerConfig.setKrb5File(krb5File); } String jaasFile = PropertiesUtils.getProperty(properties, KafkaConstants.CANAL_MQ_KAFKA_KERBEROS_JAAS_FILE); if (!StringUtils.isEmpty(jaasFile)) { kafkaProducerConfig.setJaasFile(jaasFile); } }
@Override public void stop() { try { logger.info("## stop the kafka producer"); if (producer != null) { producer.close(); } super.stop(); } catch (Throwable e) { logger.warn("##something goes wrong when stopping kafka producer:", e); } finally { logger.info("## kafka producer is down."); } }
@Override public void send(MQDestination mqDestination, Message message, Callback callback) { ExecutorTemplate template = new ExecutorTemplate(sendExecutor);
try { List result; if (!StringUtils.isEmpty(mqDestination.getDynamicTopic())) { Map<String, Message> messageMap = MQMessageUtils.messageTopics(message, mqDestination.getTopic(), mqDestination.getDynamicTopic());
for (Map.Entry<String, Message> entry : messageMap.entrySet()) { final String topicName = entry.getKey().replace('.', '_'); final Message messageSub = entry.getValue(); template.submit((Callable) () -> { try { return send(mqDestination, topicName, messageSub, mqProperties.isFlatMessage()); } catch (Exception e) { throw new RuntimeException(e); } }); }
result = template.waitForResult(); } else { result = new ArrayList(); List<Future> futures = send(mqDestination, mqDestination.getTopic(), message, mqProperties.isFlatMessage()); result.add(futures); }
producer.flush(); for (Object obj : result) { List<Future> futures = (List<Future>) obj; for (Future future : futures) { try { future.get(); } catch (InterruptedException | ExecutionException e) { throw new RuntimeException(e); } } }
callback.commit(); } catch (Throwable e) { logger.error(e.getMessage(), e); callback.rollback(); } finally { template.clear(); } }
private List<Future> send(MQDestination mqDestination, String topicName, Message message, boolean flat) { List<ProducerRecord<String, byte[]>> records = new ArrayList<>(); Integer partitionNum = MQMessageUtils.parseDynamicTopicPartition(topicName, mqDestination.getDynamicTopicPartitionNum()); if (partitionNum == null) { partitionNum = mqDestination.getPartitionsNum(); } if (!flat) { if (mqDestination.getPartitionHash() != null && !mqDestination.getPartitionHash().isEmpty()) { EntryRowData[] datas = MQMessageUtils.buildMessageData(message, buildExecutor); Message[] messages = MQMessageUtils.messagePartition(datas, message.getId(), partitionNum, mqDestination.getPartitionHash(), this.mqProperties.isDatabaseHash()); int length = messages.length; for (int i = 0; i < length; i++) { Message messagePartition = messages[i]; if (messagePartition != null) { records.add(new ProducerRecord<>(topicName, i, null, CanalMessageSerializerUtil.serializer(messagePartition, mqProperties.isFilterTransactionEntry()))); } } } else { final int partition = mqDestination.getPartition() != null ? mqDestination.getPartition() : 0; records.add(new ProducerRecord<>(topicName, partition, null, CanalMessageSerializerUtil.serializer(message, mqProperties.isFilterTransactionEntry()))); } } else { EntryRowData[] datas = MQMessageUtils.buildMessageData(message, buildExecutor); List<FlatMessage> flatMessages = MQMessageUtils.messageConverter(datas, message.getId()); for (FlatMessage flatMessage : flatMessages) { if (mqDestination.getPartitionHash() != null && !mqDestination.getPartitionHash().isEmpty()) { FlatMessage[] partitionFlatMessage = MQMessageUtils.messagePartition(flatMessage, partitionNum, mqDestination.getPartitionHash(), this.mqProperties.isDatabaseHash()); int length = partitionFlatMessage.length; for (int i = 0; i < length; i++) { FlatMessage flatMessagePart = partitionFlatMessage[i]; if (flatMessagePart != null) { records.add(new ProducerRecord<>(topicName, i, null, JSON.toJSONBytes(flatMessagePart, JSONWriter.Feature.WriteNulls, JSONWriter.Feature.LargeObject))); } } } else { final int partition = mqDestination.getPartition() != null ? mqDestination.getPartition() : 0; records.add(new ProducerRecord<>(topicName, partition, null, JSON.toJSONBytes(flatMessage, JSONWriter.Feature.WriteNulls, JSONWriter.Feature.LargeObject))); } } }
return produce(records); }
private List<Future> produce(List<ProducerRecord<String, byte[]>> records) { List<Future> futures = new ArrayList<>(); for (ProducerRecord record : records) { futures.add(producer.send(record)); }
return futures; }
}
|