我使用 kafka 作为微服务架构的消息总线,因此多个服务监听一个主题以获取消息。因此,服务高度依赖于要直播的主题。
但是,在很多情况下,我都得到了leader not available,broker not available以及leader= - 1主题。
现在,我不确定我是否可以依赖 kafka 主题,因为当主题出现问题并导致平台出现问题时,服务就会中断。
有人可以对这些主题的可靠性和可靠性有所了解吗,如果我们可以解决上述问题,我们是否可以恢复。
180718 12:43:04 [ERROR] Can't start server: Bind on TCP/IP port. Got error: 10048: Only one usage of each socket address (protocol/network address/port) is normally permitted.
180718 12:43:04 [ERROR] Do you already have another mysqld server running on port: 3306 ?
180718 12:43:04 [ERROR] Aborting
Run Code Online (Sandbox Code Playgroud)
有什么解决办法吗?Mysql 在基于 Windows 的服务器上运行。请给出最佳解决方案。。
我的主题只有一个分区,但我需要实现多重处理。我有大量异步生成的消息,我想异步读取所有这些消息并提交每条消息。
我正在运行 kafka-connect 分布式设置。
我正在使用单机/进程设置(仍处于分布式模式)进行测试,效果很好,现在我正在使用 3 个节点(和 3 个连接进程),日志不包含错误,但是当我提交 s3-connector 时通过rest-api请求,它返回:{"error_code":409,"message":"Cannot complete request because of a conflicting operation (e.g. worker rebalance)"}。
当我停止其中一个节点上的 kafka-connect 进程时,我实际上可以提交作业并且一切正常。
我的集群中有 3 个代理,主题的分区号是 32。
这是我尝试启动的连接器:
{
"name": "s3-sink-new-2",
"config": {
"connector.class": "io.confluent.connect.s3.S3SinkConnector",
"tasks.max": "32",
"topics": "rawEventsWithoutAttribution5",
"s3.region": "us-east-1",
"s3.bucket.name": "dy-raw-collection",
"s3.part.size": "64000000",
"flush.size": "10000",
"storage.class": "io.confluent.connect.s3.storage.S3Storage",
"format.class": "io.confluent.connect.s3.format.avro.AvroFormat",
"schema.generator.class": "io.confluent.connect.storage.hive.schema.DefaultSchemaGenerator",
"partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
"partition.duration.ms": "60000",
"path.format": "\'year\'=YYYY/\'month\'=MM/\'day\'=dd/\'hour\'=HH",
"locale": "US",
"timezone": "GMT",
"timestamp.extractor": "RecordField",
"timestamp.field": "procTimestamp",
"name": "s3-sink-new-2"
}
}
Run Code Online (Sandbox Code Playgroud)
日志中没有任何内容表明有问题,我真的迷失在这里。
我正在尝试删除 Kafka 主题,__Consumer_offset因为它给我的经纪人带来了很多困惑。当我这样做时,它说该主题无法标记为删除。
我正在使用 Zookeeper cli 来删除它,例如rmr /brokers/topic __consumer_offset,但它不起作用!
apache-kafka kafka-consumer-api kafka-producer-api apache-zookeeper
我正在运行 kafka 2.13-2.4.1 并配置用 java 编写的 kafka 客户端(消费者)和 kafka 集群(3 个节点,每个节点有一个代理)之间的 SSL 连接。我通过Confluence 的文档使用了官方文档,该文档具有单向身份验证(客户端没有证书),但它不起作用,所以我不得不使用两种身份验证,然后消费者和生产者控制台都通过 SSL 进行良好的通信,但是当我使用我的java消费者应用程序时:
package kafkaconsumerssl;
import java.time.Duration;
import java.util.Arrays;
import java.util.Properties;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.KafkaException;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
public class KafkaConsumerSSLTest {
public static void main(String[] args) throws KafkaException {
Properties props = new Properties();
props.put("security.protocol", "SSL");
props.put("ssl.endpoint.identification.algorithm=", "");
props.put("ssl.truststore.location","/var/private/ssl/kafka.client.truststore.jks");
props.put("ssl.truststore.password","*******");
props.put("ssl.keystore.location", "/var/private/ssl/kafka.client.keystore.jks");
props.put("ssl.keystore.password", "********");
props.put("ssl.key.password", "*******");
props.put("acks", "all");
props.put("retries", "0");
props.setProperty("zk.connnect", "172.31.32.219:2181,172.31.41.226:2181,172.31.33.133:2181");
props.setProperty("group.id", "ConsumersTest");
props.setProperty("auto.offset.reset","earliest");
props.setProperty("enable.auto.commit", "true");
props.setProperty("auto.commit.interval.ms", "1000");
props.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); …Run Code Online (Sandbox Code Playgroud) java ssl apache-kafka kafka-consumer-api apache-kafka-security
我试图了解在 kafka 消费者中处理需要更长时间处理的记录的更好选择是什么?我进行了一些测试来理解这一点,并观察到我们可以通过修改max.poll.records或来控制这一点max.poll.interval.ms。
现在我的问题是,什么是更好的选择?请建议。
最近,我开始学习使用kafka工作。我正在开发的项目使用sarama。
为了阅读消息,我使用ConsumerGroup.
foo如果返回,我需要在一段时间后再次阅读该消息false。如何才能做到这一点?
func (consumer *Consumer) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for message := range claim.Messages() {
if ok := foo(message); ok {
session.MarkMessage(message, "")
} else {
// ???
}
}
return nil
}
Run Code Online (Sandbox Code Playgroud) 我需要设置两个参数min.insync.replicas和acks。官方文档说该参数min.insync.replicas是broker的参数。我是否正确理解,对于所有主题,都应该在 server.properties 文件中指定它?其中之一是使用命令 kafka.config.sh。Acks参数只能在配置生产者时设置,例如从应用程序?更改文件 Producer.properties 没有帮助吗?
我是 Kafka 新手,想要获取每个分区的 Kafka 主题的位置。我在文档中看到 - https://kafka- python.readthedocs.io/en/master/apidoc/KafkaAdminClient.html#kafkaadminclient - 偏移量可以通过函数获得KafkaAdminClient.list_consumer_group_offsets,但我没有看到这样的方法那里的位置。
有人知道我怎样才能得到它吗?