Skip to main content

Posts

Re: can kafka state stores be used as a application level cache by application to modify it from outside the stream topology?

Thanks... i will try increasing the memory in case you don't spot anything wrong with the code. Other service also have streams and global k table but they use spring-kafka, but i think that should not matter, and it should work with normal kafka-streams code unless i am missing some configuration/setting here On Wed, May 27, 2020 at 10:26 PM Matthias J. Sax < mjsax@apache.org > wrote: > There is no hook. Only a restore listener, but this one is only used > during startup when the global store is loaded. It's not sure during > regular processing. > > Depending on your usage, maybe you can switch to a global store instead > of GlobalKTable? That way, you can implement a custom `Processor` and > add a hook manually? > > I don't see anything wrong with your setup. Unclear if/why the global > store would require a lot of memory... > > > -Matthias > > On 5/27/20 7:41 AM, Pushkar Deole wrote: > > Ma...

Re: Request for adding to contributors list

Sent previous mail w/o seeing this :| Thanks for looking up my userID and adding to the list! -Guru On 2020/05/27 16:51:41, "Matthias J. Sax" < mjsax@apache.org > wrote: > Seems Guru create the ticket. > > Added user `tguruprasad` to the list of contributors. You can know > self-assign tickets. > > > -Matthias > > On 5/26/20 11:41 PM, Luke Chen wrote: > > Hi Guruprasad, > > What Matthias was asking, is what your JIRA account username is? > > So that he or other committer can help you grant you the permission to > > assign JIRA tickets. > > If you haven't got the JIRA account yet, you can sign up here: > > https://issues.apache.org/jira/secure/Signup!default.jspa > > > > Thanks. > > Luke > > > > On Wed, May 27, 2020 at 12:32 PM Guruprasad Tahasildar < gct@cs.ucla.edu > > > wrote: > > > >> Link to Jira I would lik...

Re: Request for adding to contributors list

My bad, here is my Jira ID: tguruprasad Thanks, On 2020/05/27 06:41:02, Luke Chen < showuon@gmail.com > wrote: > Hi Guruprasad, > What Matthias was asking, is what your JIRA account username is? > So that he or other committer can help you grant you the permission to > assign JIRA tickets. > If you haven't got the JIRA account yet, you can sign up here: > https://issues.apache.org/jira/secure/Signup!default.jspa > > Thanks. > Luke > > On Wed, May 27, 2020 at 12:32 PM Guruprasad Tahasildar < gct@cs.ucla.edu > > wrote: > > > Link to Jira I would like to address: > > https://issues.apache.org/jira/browse/KAFKA-10047 > > > > Thanks, > > > > On 2020/05/26 17:26:06, "Matthias J. Sax" < mjsax@apache.org > wrote: > > > What is your Jira ID? > > > > > > -Matthias > > > > > > On 5/24/20 8:31 PM, Guruprasad Taha...

Re: can kafka state stores be used as a application level cache by application to modify it from outside the stream topology?

-----BEGIN PGP SIGNATURE----- iQIzBAEBCAAdFiEEI8mthP+5zxXZZdDSO4miYXKq/OgFAl7Om6sACgkQO4miYXKq /Oj3IQ//UE4oVrq9+H2qO1CjJlscx85m7t3UVgU2ERdYWGqcevOjoxuu/L42AGtB lZTRfvFLVRyrqxUXBMmneffVOnXX3BqKIfHVzAzHoqW08+hxTXB+4RSAesMSQsN0 vuVneNNqGiV4H9TFenk60BVQzXHE5HfSYfN9U0jcQqJfE2xK6bEhz3IHOWg8LIhD o6wDx5VpQY3yWMU+oamujE7bh4XheUgIF/+YuIzpYnQ/42VYRNyGGJkTM9Zvzm78 Ziu7W7uIYIlXew0rEbldBidBAXjKFzZ+8MBislL9YhbbLGSkgEchqV7cZXApW5cX slOxg0gy2ZJ+LdTMTFXyF5yi65T35ZhdMfY/L99yZfXzAbjxKcf5xuAX5pNrVrxl PPUCb/hjrIA+qa8OWcMcsBd4wgHBKzjRE+Y3nb9uoAouGB3XopjO0P9N5pfBin+P p+xjK74RqWny9KAGBpQEV2hZTYftywO5JiyqkX0KMGD7QYS1relY2IBM255XS4QT f2y9YYoNR87qFRROYX+V7iZCvQ1iQsc0c0BY9rtJGZhoVj+Q1BtRy8Q4sdKT+qM4 eWSgZS15NGIN8J6ByeMJEg5buWOSf3opxPLCRj5tGehH8B7GxpgOY4oE1tyjwfm+ KkgQ9iy58gn5IbKuXnPUCoVEyvKel5+gegZx/l3jeLlIvtWZO5k= =cyKC -----END PGP SIGNATURE----- There is no hook. Only a restore listener, but this one is only used during startup when the global store is loaded. It's not sure during regular...

Re: Request for adding to contributors list

-----BEGIN PGP SIGNATURE----- iQIzBAEBCAAdFiEEI8mthP+5zxXZZdDSO4miYXKq/OgFAl7Omp0ACgkQO4miYXKq /Oizxw//f5t2OSb0kzaJWdqaB7eKmiotKIl6H+fJM8hNY/REpayYRBQTx/+qhxEO vf3W+EY/2Wmw8wVpv0pbW0k6IoU1HJe1ukCeaCtjjp+MGQLKhq+p07J5Kfb75qsq vFwrfavFELIwkmN2wt2U7fVpWYGTJ+nnwy4+ZvsGN78XhDNQMwuBKwiwk9Lv6eKC 2zXkIoTuN5mKpz+c3njywoBGyGnuF76LjOL8SGZXHu+IIgvVakVYT0KI1Csr50Yc NQC1G5PyzKA10kLZHW1sJM1aZ4ftFFtV/dJk8OckvbAkmJowE/7Wc2RSK8KYrIfE NxVfwb/IVAPU+gY4X3GdtXDKa4VmnGMy7JYpgkfe3DXEsg6t+eFWG/EP0fG2lWyA 2PHlo7w3WevKQzZZ4w+3nHlY2gqpjgmWI24UX57WRKuxrk2bRrt+kXJCA3Vr/BDB M4KcbiKkeFAP+y1e18PguWOXW7+3njiuLoAzozCODmOAM0ivHDyWrpu/ddIQ/kb1 Ij2sYUuwVldf/5vwpBB0d735CX3fYNzORQs5bdKH3q5Tl12PaeTPNpiWkwHX5zEA ydjwuuxk4WOdsIqPs0IreYAnzAatfQZGtvIZyNiVkCo8CbKh2h7qWuQrZ+tWkk+B NTsFXOagW9xGSAZUnVQUfBbRUJpCXcc/SXdHkfrfHpLzSVoVI18= =8Hh8 -----END PGP SIGNATURE----- Seems Guru create the ticket. Added user `tguruprasad` to the list of contributors. You can know self-assign tickets. -Matthias On 5/26/20 ...

Re: Kafka retention policy per topic not working as expected

Thanks Christopher, Hrm, I checked the logs and retention.ms has been enabled somewhere around October last year. {"log":"[2019-10-20 03:01:41,818] INFO Processing override for entityPath: topics/ETD-TEST with config: Map( retention.ms -\u003e 1210000000) (kafka.server.DynamicConfigManager)\n","stream":"stdout","time":"2019-10-20T03:01:41.818402555Z"} So don't think the first hypothesis of it being turned on recently sticks. I'll have to investigate the other possibility you mentioned but we only have 46 topics with roughly 4 million events total. I'll continue to dig but thank you for your thoughts and chiming in! On Sun, May 24, 2020 at 8:10 PM Christopher Smith < cbsmith@gmail.com > wrote: > It's a bit weird having the default policy for your brokers being compact, > but yes, the policy for the topic overrides the broker policy. > > What you are seeing on...

Re: can kafka state stores be used as a application level cache by application to modify it from outside the stream topology?

Matthias, I tried with default store as well but getting same error, can you please check if I am initializing the global store in the right way: public void setupGlobalCacheTables(String theKafkaServers) { Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, DEFAULT_APPLICATION_ID); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, theKafkaServers); StreamsBuilder streamsBuilder = new StreamsBuilder(); groupCacheTable = streamsBuilder.globalTable(GROUP_CACHE_TOPIC, Consumed.with(Serdes.String(), GroupCacheSerdes.groupCache()), Materialized.as(GROUP_CACHE_STORE_NAME)); Topology groupCacheTopology = streamsBuilder.build(); kafkaStreams = new KafkaStreams(groupCacheTopology, props); kafkaStreams.start(); Runtime.getRuntime().addShutdownHook(new Thread(() -> { LOG.info("Stopping the stream"); kafkaStreams.close(); })); } On Wed, May 27, 2020 at 5:06 PM Pushkar Deole < pd...