标签: confluent-platform

**Kafka** 跨区域数据中心之间的双向同步

我有一个部署,我们使用 kafka 从服务发送消息。但是我们需要在所有地区都拥有 Kafka 大师。因此,一旦消息在 1 个数据中心推送,它应该在其他数据中心同步。当它再次在其他数据中心完成时,它应该同步回来。Mirror Maker 可以提供从 1 到其他的同步,但我如何实现双向同步?

apache-kafka confluent-platform

4
推荐指数
1
解决办法
2382
查看次数

Confluence Kafka:消费者不会从头开始读取主题中的所有分区

我有一个有 40 个分区的主题。设置是这样的:

def on_assign (c,ps):
    for p in ps:
        p.offset=0
    print ps
    c.assign(ps)

conf = {'bootstrap.servers': 'localhost:9092'
        'enable.auto.commit' : False,
        'group.id' : 'confluent_consumer',
        'default.topic.config': {'auto.offset.reset': 'earliest'}
        }
consumer = Consumer(**conf)
consumer.subscribe(['topic.source'], on_assign=on_assign)

msg = consumer.poll(timeout=100000)
print "Topic is %s: | Partition is %d: | Offset is : %d | key is :%s " % (msg.topic(), msg.partition(), msg.offset(), msg.key())
Run Code Online (Sandbox Code Playgroud)

我想从偏移量 0 读取主题的所有分区topic.source。但我没有看到所有分区都会发生这种情况。对于某些分区,它从特定的偏移量读取,我假设这是提交的偏移量,group.id每次更改也没有帮助。如何从头开始读取该主题的所有分区,而不考虑提交的偏移量?

我打印ps出来on_assign(),它为所有 40 个分区打印了这样的内容:

[TopicPartition{topic=topic.source,partition=0,offset=0,error=None},TopicPartition{topic=topic.source,partition=1,offset=0,error=None}....] and so on
Run Code Online (Sandbox Code Playgroud)

python apache-kafka kafka-consumer-api confluent-platform

4
推荐指数
1
解决办法
8460
查看次数

org.apache.kafka.common.KafkaException:io.confluence.kafka.serializers.KafkaAvroSerializer

尝试使用以下代码构建一个kafka消费者

import org.apache.kafka.clients.consumer.ConsumerConfig;
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.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
// set up consumer
final Properties consumerProps = new Properties();
consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, CLUSTER.bootstrapServers());
consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "consumer-tutorial");
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
        io.confluent.kafka.serializers.KafkaAvroSerializer.class);
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
        io.confluent.kafka.serializers.KafkaAvroSerializer.class);
// transactional API
consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
// consumer --from-beginning
consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
consumerProps.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "10000");
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
consumerProps.put("zookeeper.connect", CLUSTER.zookeeperConnect());
consumerProps.put("schema.registry.url", CLUSTER.schemaRegistryUrl());
final KafkaConsumer<GenericRecord, GenericRecord> consumer = new KafkaConsumer<GenericRecord, GenericRecord>(consumerProps);
consumer.subscribe(Collections.singletonList(inputTopic));
Run Code Online (Sandbox Code Playgroud)

但因错误而失败

org.apache.kafka.common.KafkaException: Failed to construct kafka consumer
    at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:765)
    at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:633)
    at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:615)
    at com.telefonica.app.test_consumer.KafkaETLConsumerTest.testRunConsumer(KafkaETLConsumerTest.java:192)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) …
Run Code Online (Sandbox Code Playgroud)

apache-kafka confluent-platform

4
推荐指数
1
解决办法
7315
查看次数

Windows 中的 Confluence 平台

是否可以在 Windows 上运行完整的 Confluence 平台?如果不是,运行 Confluence 平台的最佳方式是什么?

谢谢

windows apache-kafka confluent-platform

4
推荐指数
1
解决办法
2239
查看次数

ElasticsearchSinkConnector 无法将数据反序列化到 Avro

我创建了最简单的 kafka 接收器连接器配置,我使用的是 confluence 4.1.0:

{
  "connector.class": 
  "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
  "type.name": "test-type",
  "tasks.max": "1",
  "topics": "dialogs",
  "name": "elasticsearch-sink",
  "key.ignore": "true",
  "connection.url": "http://localhost:9200",
  "schema.ignore": "true"
}
Run Code Online (Sandbox Code Playgroud)

在主题中,我将消息保存为JSON

{ "topics": "resd"}
Run Code Online (Sandbox Code Playgroud)

但在结果中我得到一个错误:

引起:org.apache.kafka.common.errors.SerializationException:反序列化 id -1 的 Avro 消息时出错 引起:org.apache.kafka.common.errors.SerializationException:未知的魔术字节!

apache-kafka apache-kafka-connect confluent-platform

4
推荐指数
1
解决办法
2254
查看次数

我们如何强制融合 kafka 连接 s3 接收器进行冲洗

我设置了 kafka 连接 s3 接收器,持续时间设置为 1 小时,并且我设置了一个相当大的刷新计数,比如 10,000。现在如果kafka通道中的消息不多,s3 sink会尝试将它们缓存在内存中,等待它累积到flush计数,然后将它们一起上传并将偏移量提交给自己的消费者组。

但是想想这种情况。如果在频道里,我只发送5000条消息。然后没有 s3 水槽冲洗。然后时间长了,这5000条消息最终会因为保留时间的原因从kafka中驱逐出去。但是这些消息仍然在 s3 sink 的内存中,而不是在 s3 中。这是非常危险的,例如,如果我们重新启动 s3 sink 或运行 s3 sink 的机器就崩溃了。然后我们丢失了那 5,000 条消息。我们无法从 kafka 中再次找到它们,因为它已经被删除了。

这会发生在 s3 sink 上吗?或者有一些设置会强制它在一段时间后刷新?

amazon-s3 apache-kafka apache-kafka-connect confluent-platform

4
推荐指数
1
解决办法
1363
查看次数

Kafka 与 Confluent Kubernetes Helm Charts = Schema Registry WakeupException

我的主要问题:为什么模式注册表会崩溃?

外围问题:如果我为每个 zookeeper/kafka/schema-registry 配置了一个服务器,为什么每个 pod 都启动?其他一切看起来基本正确吗?

?  helm repo update
<snip>

?  helm install --values values.yaml --name my-confluent-oss confluentinc/cp-helm-charts
<snip>

?  helm list
NAME                REVISION    UPDATED                     STATUS      CHART                   APP VERSION NAMESPACE
my-confluent-oss    1           Sat Oct 20 19:09:08 2018    DEPLOYED    cp-helm-charts-0.1.0    1.0         default  

?  kubectl get pods
NAME                                                   READY     STATUS             RESTARTS   AGE
my-confluent-oss-cp-kafka-0                            2/2       Running            0          20m
my-confluent-oss-cp-schema-registry-59d8877584-c2jc7   1/2       CrashLoopBackOff   7          20m
my-confluent-oss-cp-zookeeper-0                        2/2       Running            0          20m
Run Code Online (Sandbox Code Playgroud)

values.yaml的如下。我已经用helm install --debug --dry-run. 我只是禁用持久性,设置单个服务器(这是在 VM 中运行的开发设置),并暂时禁用额外服务,直到我获得基础工作:

cp-kafka:
  brokers: 1 …
Run Code Online (Sandbox Code Playgroud)

apache-kafka kubernetes kubernetes-helm confluent-schema-registry confluent-platform

4
推荐指数
1
解决办法
2421
查看次数

kafka-connect 在分布式模式下返回 409

我正在运行 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)

日志中没有任何内容表明有问题,我真的迷失在这里。

apache-kafka apache-kafka-connect confluent-platform

4
推荐指数
1
解决办法
1936
查看次数

如何将基本身份验证传递给 Confluent Schema Registry?

我想从融合的云主题读取数据,然后写入另一个主题。

在本地主机上,我没有遇到任何重大问题。但是confluent cloud的schema registry需要传递一些我不知道如何输入的身份验证数据:

basic.auth.credentials.source=USER_INFO

schema.registry.basic.auth.user.info=:

schema.registry.url= https://xxxxxxxxxx.confluent.cloudBlockquote

以下是当前代码:

import com.databricks.spark.avro.SchemaConverters
import io.confluent.kafka.schemaregistry.client.{CachedSchemaRegistryClient, SchemaRegistryClient}
import io.confluent.kafka.serializers.AbstractKafkaAvroDeserializer
import org.apache.avro.Schema
import org.apache.avro.generic.GenericRecord
import org.apache.spark.sql.SparkSession

object AvroConsumer {
  private val topic = "transactions"
  private val kafkaUrl = "http://localhost:9092"
  private val schemaRegistryUrl = "http://localhost:8081"

  private val schemaRegistryClient = new CachedSchemaRegistryClient(schemaRegistryUrl, 128)
  private val kafkaAvroDeserializer = new AvroDeserializer(schemaRegistryClient)

  private val avroSchema = schemaRegistryClient.getLatestSchemaMetadata(topic + "-value").getSchema
  private var sparkSchema = SchemaConverters.toSqlType(new Schema.Parser().parse(avroSchema))

  def main(args: Array[String]): Unit = {
    val spark = SparkSession
      .builder
      .appName("ConfluentConsumer")
      .master("local[*]")
      .getOrCreate() …
Run Code Online (Sandbox Code Playgroud)

apache-spark databricks confluent-schema-registry spark-structured-streaming confluent-platform

4
推荐指数
1
解决办法
2827
查看次数

Kafka 依赖项 - ccs 与 ce

要开发我的 Kafka 连接器,我需要添加一个 connect-API 依赖项。

我应该使用哪一种?

例如 mongodb 连接器使用来自maven central 的connect-api

但是来自开发指南的链接转到https://packages.confluent.io/maven/org/apache/kafka/connect-api/5.5.0-ccs/旁边5.5.0-ccs还有5.5.0-ce版本。

所以,此时最后的版本是:

所有三个变体之间有什么区别?

我应该使用哪一种?

java maven apache-kafka apache-kafka-connect confluent-platform

4
推荐指数
1
解决办法
1110
查看次数