apache / apache/pinot

Need Help - Realtime table | External Kafka Topic Configuration with SSL

Open
#8,041 8 comments 0 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
6.1k
Forks
1.5k
Avg merge
1d 21h
Merged PRs (30d)
189

Description

Hi Team,

I am trying to connect my existing Kafka topic with a real-time table but it's showing a timeout issue. can you please help me configure the real-time table with SSL?

All of my configuration details are mentioned below.

@moszutij @mayankshriv

### ERROR

```
2022/01/19 10:22:07.249 INFO [AddTableCommand] [main] bash pinot-admin.sh AddTable \
-schemaFile ./../../../app/empx_schema_config.json \
-tableConfigFile ./../../../app//empx_table_config.json \
-controllerHost localhost \
-controllerPort 9000 \
-exec > > > > >
[0.001s][warning][gc] -Xloggc is deprecated. Will use -Xlog:gc:gc-pinot-controller.log instead.
SLF4J: Class path contains multiple SLF4J bindings.
SLF4J: Found binding in [jar:file:/opt/pinot/lib/pinot-all-0.9.3-jar-with-dependencies.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: Found binding in [jar:file:/opt/pinot/plugins/pinot-file-system/pinot-s3/pinot-s3-0.9.3-shaded.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: Found binding in [jar:file:/opt/pinot/plugins/pinot-metrics/pinot-yammer/pinot-yammer-0.9.3-shaded.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: Found binding in [jar:file:/opt/pinot/plugins/pinot-metrics/pinot-dropwizard/pinot-dropwizard-0.9.3-shaded.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: Found binding in [jar:file:/opt/pinot/plugins/pinot-environment/pinot-azure/pinot-azure-0.9.3-shaded.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: Found binding in [jar:file:/opt/pinot/plugins/pinot-input-format/pinot-parquet/pinot-parquet-0.9.3-shaded.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: See http://www.slf4j.org/codes.html#multiple_bindings for an explanation.
SLF4J: Actual binding is of type [org.apache.logging.slf4j.Log4jLoggerFactory]
WARNING: sun.reflect.Reflection.getCallerClass is not supported. This will impact performance.
WARNING: An illegal reflective access operation has occurred
WARNING: Illegal reflective access by org.codehaus.groovy.reflection.CachedClass (file:/opt/pinot/lib/pinot-all-0.9.3-jar-with-dependencies.jar) to method java.lang.Object.finalize()
WARNING: Please consider reporting this to the maintainers of org.codehaus.groovy.reflection.CachedClass
WARNING: Use --illegal-access=warn to enable warnings of further illegal reflective access operations
WARNING: All illegal access operations will be denied in a future release
2022/01/19 10:22:01.432 INFO [AddTableCommand] [main] Executing command: AddTable -tableConfigFile ./../../../app//blx0583t_table_config.json -schemaFile ./../../../app/blx0583t_schema_config.json -controllerProtocol http -controllerHost localhost -controllerPort 9000 -user null -password [hidden] -exec
2022/01/19 10:22:07.249 INFO [AddTableCommand] [main] {"_code":500,"_error":"org.apache.kafka.common.errors.TimeoutException: Timeout expired while fetching topic metadata"}

```

### How I am consuming my topic in console

```
./confluent/bin/kafka-avro-console-consumer \
--bootstrap-server xx-xx-xx.com:443 \
--topic employee-data-feeds \
--property schema.registry.url=http://yy-yyyy.com/ \
--consumer.config client-ssl.properties \
--from-beginning \
--property print.key=true \
--max-messages 10

//client-ssl.properties FILE CODE

bootstrap.servers= xx-xx-xx.com:443
schema.registry.url==http://yy-yyyy.com/
security.protocol=SSL
sasl.mechanism=PLAIN
ssl.keystore.location=./kafka-cert/keystore.jks
ssl.keystore.password=*******
ssl.truststore.location=./kafka-cert/truststore.jks
ssl.truststore.password=*******
ssl.key.password=*******
ssl.client.auth=required
ssl.endpoint.identification.algorithm=

```
### My Topic Avro output

```
{
"CERT":{"string":"XXXXXXXX"},
"SEQ":{"string":"0"},
"NAME":{"string":"DANA JHON"},
"UPDT_TMSTM":{"long":1598483669419893},
"beforeData":null,
"headers":{
"operation":"REFRESH",
"changeSequence":"","timestamp":"",
"streamPosition":"",
"transactionId":"",
"changeMask":null,
"columnMask":null,
"transactionEventCounter":null,
"transactionLastEvent":null
}
}
```

### My Pinot Schema (empx_schema_config.json)

```
{
"schemaName": "empx_schema",
"dimensionFieldSpecs": [
{
"name": "CERT",
"dataType": "STRING"
},
{
"name": "SEQ",
"dataType": "STRING"
},
{
"name": "NAME",
"dataType": "STRING"
}
],
"metricFieldSpecs": [
{
"name": "NBR",
"dataType": "INT"
}
],
"dateTimeFieldSpecs": [
{
"name": "UPDT_TMSTM",
"dataType": "LONG",
"format": "1:MILLISECONDS:EPOCH",
"granularity": "1:MILLISECONDS"
}
]
}

```

### My Pinot Realtime Table Config (empx_table_config.json)

{
"tableName": "empx_table",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "UPDT_TMSTM",
"timeType": "MILLISECONDS",
"schemaName": "empx_schema",
"replication": "1",
"replicasPerPartition": "1",
"retentionTimeUnit": null,
"retentionTimeValue": null,
"completionConfig": null,
"crypterClassName": null,
"peerSegmentDownloadScheme": null
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP",
"invertedIndexColumns": [],
"createInvertedIndexDuringSegmentGeneration": false,
"rangeIndexColumns": [],
"sortedColumn": [],
"bloomFilterColumns": [],
"bloomFilterConfigs": null,
"noDictionaryColumns": [],
"onHeapDictionaryColumns": [],
"varLengthDictionaryColumns": [],
"enableDefaultStarTree": false,
"starTreeIndexConfigs": null,
"enableDynamicStarTreeCreation": false,
"segmentPartitionConfig": null,
"columnMinMaxValueGeneratorMode": null,
"aggregateMetrics": false,
"nullHandlingEnabled": false,
"streamConfigs": {
"streamType": "kafka",
"stream.kafka.consumer.type": "LowLevel",
"stream.kafka.topic.name": "employee-data-feeds",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.avro.KafkaAvroMessageDecoder",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory",
"stream.kafka.zk.broker.url": "localhost:2191/kafka",
"stream.kafka.broker.list": "localhost:19092",
"stream.kafka.hlc.zk.connect.string": "localhost:2191/kafka",
"stream.kafka.hlc.bootstrap.server": "xx-xx-xx.com:443",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest",
"stream.kafka.consumer.timeout":"",
"schema.registry.url": "http://yy-yyyy.com/",
"security.protocol": "SSL",
"ssl.truststore.location": "./../../../app/umr-certs/truststore.jks",
"ssl.keystore.location": "./../../../app/umr-certs/keystore.jks",
"ssl.truststore.password": "temppassword***",
"ssl.keystore.password": "temppassword***",
"ssl.key.password": "temppassword***",
"realtime.segment.flush.threshold.time": "24h",
"realtime.segment.flush.threshold.size": "100M",
"realtime.segment.flush.threshold.rows": "100",
"realtime.segment.flush.segment.size": "100M"
}
},
"metadata": {
"customConfigs": {}
}
}

Contributor guide

Open the contributing guide

Research direction

Start with the provided empx_table_config.json streamConfigs and the org.apache.pinot.plugin.stream.kafka20.KafkaConsumerFactory entry point, then compare them with client-ssl.properties and the AddTable timeout output. Reproduce the external Kafka SSL setup and trace topic metadata retrieval; done means the configuration is validated or the required code and tests are clearly identified.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, kafka
Domain
databases, stream-processing
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
20/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.