# Streaming (Kafka, Spark, Flink)

**URL:** https://community.neo4j.com/c/integrations/stream-processing/22.md

[Latest](https://community.neo4j.com/latest.md) · [Categories](https://community.neo4j.com/categories.md) · [Tags](https://community.neo4j.com/tags.md)

---

## [About the Streaming (Kafka, Spark, Flink) category](https://community.neo4j.com/t/about-the-streaming-kafka-spark-flink-category/33)

<div class="topic-metadata">

**Author:** [@michael.hunger](https://community.neo4j.com/u/michael.hunger)\
**Replies:** 0\
**Last updated:** [June 11, 2018, 3:54pm UTC](https://community.neo4j.com/t/about-the-streaming-kafka-spark-flink-category/33 "2018-06-11T15:54:19Z")

</div>

Stream processing is used for both service integration and large scale data processing in Apache Kafka, Spark and, Flink. Neo4j integrates with all these libraries. If you run into a bug, rather than posting it here, pl…

---

## [📝Blog: Real-Time Supply Chain Event Streaming With Kafka and Neo4j](https://community.neo4j.com/t/blog-real-time-supply-chain-event-streaming-with-kafka-and-neo4j/80705)

<div class="topic-metadata">

**Author:** [@Ari\_Neo4j](https://community.neo4j.com/u/Ari_Neo4j)\
**Replies:** 0\
**Last updated:** [August 18, 2026, 11:36pm UTC](https://community.neo4j.com/t/blog-real-time-supply-chain-event-streaming-with-kafka-and-neo4j/80705 "2026-08-18T23:36:15Z")

</div>

Real-Time Supply Chain Event Streaming With Kafka and Neo4j A Kafka producer publishes shipment events, a Python consumer writes them into Neo4j, and a live Plotly dashboard shows network health updating as events arrive…

---

## [Sink to Neo4j from Kafka filing](https://community.neo4j.com/t/sink-to-neo4j-from-kafka-filing/75779)

<div class="topic-metadata">

**Author:** [@georgelza](https://community.neo4j.com/u/georgelza)\
**Replies:** 13\
**Last updated:** [October 30, 2025, 5:29pm UTC](https://community.neo4j.com/t/sink-to-neo4j-from-kafka-filing/75779 "2025-10-30T17:29:30Z")

</div>

hi hi all, ok, this is slightly strange, posted a couple of other messages, including on the discord server. I’m going to post working examples and then the failing example. The 2 inbound topics are adults and childre…

---

## [Sink strategy not Assigned](https://community.neo4j.com/t/sink-strategy-not-assigned/73566)

<div class="topic-metadata">

**Author:** [@adityapesu](https://community.neo4j.com/u/adityapesu)\
**Replies:** 3\
**Last updated:** [July 21, 2025, 4:39pm UTC](https://community.neo4j.com/t/sink-strategy-not-assigned/73566 "2025-07-21T16:39:51Z")

</div>

hi this is aditya I am new to neo4j and kafka I am trying to load a graph into neo4j using my kafka. The problem when i try to query the rest api to load the data the task runs but fails , with the following issue. org…

---

## [Neo4j kafka sink plugin](https://community.neo4j.com/t/neo4j-kafka-sink-plugin/37819)

<div class="topic-metadata">

**Author:** [@info15](https://community.neo4j.com/u/info15)\
**Replies:** 1\
**Last updated:** [April 20, 2025, 7:10pm UTC](https://community.neo4j.com/t/neo4j-kafka-sink-plugin/37819 "2025-04-20T19:10:30Z")

</div>

Hi all, i'm trying a super simple test integration between kafka and neo4j using neo4j sink plugin but simply not work. This are my test configuration: kafka.bootstrap.servers=127.0.0.1:9092 kafka.auto.offset.reset=ear…

---

## [Neo4j streams cypher template with dynamic label](https://community.neo4j.com/t/neo4j-streams-cypher-template-with-dynamic-label/66799)

<div class="topic-metadata">

**Author:** [@ajalali](https://community.neo4j.com/u/ajalali)\
**Replies:** 3\
**Last updated:** [March 13, 2024, 10:00am UTC](https://community.neo4j.com/t/neo4j-streams-cypher-template-with-dynamic-label/66799 "2024-03-13T10:00:53Z")

</div>

Is there a way to dynamically fetch node/relationship label from events in cypher template? Preferably not using apoc. Something like: WITH event.label as nodeLabel MERGE (n:nodeLable {id: event.id}) ... The above que…

---

## [CDC AWS EventBridge integration](https://community.neo4j.com/t/cdc-aws-eventbridge-integration/65950)

<div class="topic-metadata">

**Author:** [@niqo01](https://community.neo4j.com/u/niqo01)\
**Replies:** 0\
**Last updated:** [January 24, 2024, 7:48pm UTC](https://community.neo4j.com/t/cdc-aws-eventbridge-integration/65950 "2024-01-24T19:48:12Z")

</div>

Hello, I was wondering if Neo4j could consider integrating CDC with Aws EventBridge. I am assuming this would not be drastically different than CDC and Kafka. thank you, Nicolas

---

## [Kafka sink connector delete operation](https://community.neo4j.com/t/kafka-sink-connector-delete-operation/65408)

<div class="topic-metadata">

**Author:** [@alessandra.nunez](https://community.neo4j.com/u/alessandra.nunez)\
**Replies:** 0\
**Last updated:** [December 18, 2023, 4:44pm UTC](https://community.neo4j.com/t/kafka-sink-connector-delete-operation/65408 "2023-12-18T16:44:58Z")

</div>

Hello, I need to capture changes from a MySQL database and stream it to a Neo4j db through Kafka. I already have the CDC connector for MySQL but I'm having trouble with the Neo4j sink connector. I'm using the recommende…

---

## [An apoc.periodic.iterateSpark that offloads batch jobs to Apache Spark](https://community.neo4j.com/t/an-apoc-periodic-iteratespark-that-offloads-batch-jobs-to-apache-spark/64380)

<div class="topic-metadata">

**Author:** [@neo4joe](https://community.neo4j.com/u/neo4joe)\
**Replies:** 1\
**Last updated:** [October 5, 2023, 12:58pm UTC](https://community.neo4j.com/t/an-apoc-periodic-iteratespark-that-offloads-batch-jobs-to-apache-spark/64380 "2023-10-05T12:58:38Z")

</div>

All integrations I have seen between Spark and Neo4j put Spark in the driver's seat, where Spark queries Neo4j, transforms the data it receives, and possibly sends it back to Neo4j. I would prefer to have the experience…

---

## [Push Node and relation of Neo4j 5 to Kafka topic](https://community.neo4j.com/t/push-node-and-relation-of-neo4j-5-to-kafka-topic/63363)

<div class="topic-metadata">

**Author:** [@rahul.bisen](https://community.neo4j.com/u/rahul.bisen)\
**Replies:** 2\
**Last updated:** [August 29, 2023, 1:46pm UTC](https://community.neo4j.com/t/push-node-and-relation-of-neo4j-5-to-kafka-topic/63363 "2023-08-29T13:46:27Z")

</div>

I want to push nodes and relations of neo4j 5 to kafka topic using Neo4jSourceConnector, but failing to get relation and operation based (Create and Update) nodes to kafka topic. I can able to send node data based on l…

---

## [Nodes with PointValue property cannot be successfully sink to another Neo4j instance](https://community.neo4j.com/t/nodes-with-pointvalue-property-cannot-be-successfully-sink-to-another-neo4j-instance/62731)

<div class="topic-metadata">

**Author:** [@yhhongyang](https://community.neo4j.com/u/yhhongyang)\
**Replies:** 0\
**Last updated:** [May 30, 2023, 8:28am UTC](https://community.neo4j.com/t/nodes-with-pointvalue-property-cannot-be-successfully-sink-to-another-neo4j-instance/62731 "2023-05-30T08:28:39Z")

</div>

Hi team, I have 2 questions. First I'm using neo4j streams to sync data from one neo4j instance to another, and I found that nodes with PointValue property cannot be successfully sink , it will throw error like ErrorDa…

---

## [Kafka Neo4j Connector Source, how to capture deletes?](https://community.neo4j.com/t/kafka-neo4j-connector-source-how-to-capture-deletes/61789)

<div class="topic-metadata">

**Author:** [@dimitri1](https://community.neo4j.com/u/dimitri1)\
**Replies:** 2\
**Last updated:** [April 13, 2023, 1:31pm UTC](https://community.neo4j.com/t/kafka-neo4j-connector-source-how-to-capture-deletes/61789 "2023-04-13T13:31:13Z")

</div>

Hi, I need to capture all the changes in the neo4j database and stream it through Kafka. Neo4j Stream is deprecated and the documentation refers to the Neo4j Connector. Now instead of changes being being pushed, Kafka n…

---

## [Neo4j sink using CDC source id strategy update not working](https://community.neo4j.com/t/neo4j-sink-using-cdc-source-id-strategy-update-not-working/54651)

<div class="topic-metadata">

**Author:** [@ahmed.rabei](https://community.neo4j.com/u/ahmed.rabei)\
**Replies:** 0\
**Last updated:** [April 10, 2022, 5:25pm UTC](https://community.neo4j.com/t/neo4j-sink-using-cdc-source-id-strategy-update-not-working/54651 "2022-04-10T17:25:48Z")

</div>

Hi all, Need your support. I am trying to use cdc with source id strategy and it's working fine with create and delete but not working with update. Update is working as a new node without updating old one with new dec…

---

## [Kafka Integration: How to configure streams.source.topic.nodes and streams.source.topic.relationships to capture all nodes and relationships?](https://community.neo4j.com/t/kafka-integration-how-to-configure-streams-source-topic-nodes-and-streams-source-topic-relationships-to-capture-all-nodes-and-relationships/15382)

<div class="topic-metadata">

**Author:** [@hunter](https://community.neo4j.com/u/hunter)\
**Replies:** 13\
**Last updated:** [March 8, 2022, 12:22pm UTC](https://community.neo4j.com/t/kafka-integration-how-to-configure-streams-source-topic-nodes-and-streams-source-topic-relationships-to-capture-all-nodes-and-relationships/15382 "2022-03-08T12:22:27Z")

</div>

We are using neo4j 3.5.4 version. We have nearly 200 nodes and relationships. We want to capture all nodes and relationships data even for all CURD operations through Kafka . In neo4j .conf how I need to configure fo…

---

## [Using streams.publish procedure is there a way to publish events with key in Kafka topic?](https://community.neo4j.com/t/using-streams-publish-procedure-is-there-a-way-to-publish-events-with-key-in-kafka-topic/12519)

<div class="topic-metadata">

**Author:** [@kavitakjava](https://community.neo4j.com/u/kavitakjava)\
**Replies:** 2\
**Last updated:** [February 24, 2022, 3:25am UTC](https://community.neo4j.com/t/using-streams-publish-procedure-is-there-a-way-to-publish-events-with-key-in-kafka-topic/12519 "2022-02-24T03:25:22Z")

</div>

As we know we can publish events with key in Kafka topic e.g. kafka-console-producer --topic key-value-topic --broker-list localhost:9092 --property "parse.key=true" --property "key.separator=:" key1:{value1} key2:{…

---

## [Error During Delete](https://community.neo4j.com/t/error-during-delete/46425)

<div class="topic-metadata">

**Author:** [@psfurlong](https://community.neo4j.com/u/psfurlong)\
**Replies:** 0\
**Last updated:** [October 25, 2021, 7:13pm UTC](https://community.neo4j.com/t/error-during-delete/46425 "2021-10-25T19:13:06Z")

</div>

Hi. Hoping to get a little wisdom here as to whether this is a bug or something we've configured incorrectly. Using Neo4j 3.5 and sending to a Kafka topic using the streams capability. All seems to work well for creating…

---

## [Consumer config doesn't seem to work](https://community.neo4j.com/t/consumer-config-doesnt-seem-to-work/43206)

<div class="topic-metadata">

**Author:** [@pselvakumarmca](https://community.neo4j.com/u/pselvakumarmca)\
**Replies:** 1\
**Last updated:** [August 30, 2021, 8:30pm UTC](https://community.neo4j.com/t/consumer-config-doesnt-seem-to-work/43206 "2021-08-30T20:30:41Z")

</div>

I have a sink connector which consumes messages from KAFKA, I've configured poll records as 1000 by adding this property. "kafka.max.poll.records":1000 I just followed the documentation here. But It always polls 500 …

---

## [How to configure kafka neo4j stream connector with neo4j aura](https://community.neo4j.com/t/how-to-configure-kafka-neo4j-stream-connector-with-neo4j-aura/42762)

<div class="topic-metadata">

**Author:** [@pamnanilearning](https://community.neo4j.com/u/pamnanilearning)\
**Replies:** 1\
**Last updated:** [August 19, 2021, 9:46pm UTC](https://community.neo4j.com/t/how-to-configure-kafka-neo4j-stream-connector-with-neo4j-aura/42762 "2021-08-19T21:46:39Z")

</div>

As per documentation, to use neo4j kafka stream connector, neo4j.conf is to be configure but how/where to do the same in case of neo4j aura. Moreover, if we want to used kafka connector instead of neo4j extension then h…

---

## [Neo4j Kafka Source Plugin: Re- publish everything without losing data](https://community.neo4j.com/t/neo4j-kafka-source-plugin-re-publish-everything-without-losing-data/43034)

<div class="topic-metadata">

**Author:** [@average\_legend](https://community.neo4j.com/u/average_legend)\
**Replies:** 1\
**Last updated:** [August 19, 2021, 9:10pm UTC](https://community.neo4j.com/t/neo4j-kafka-source-plugin-re-publish-everything-without-losing-data/43034 "2021-08-19T21:10:38Z")

</div>

Hi. I am using Neo4j 4.2.3 together with the neo4j-streams-4.0.8 Plugin to stream Neo4j CDC events to a kafka topic. Works okay so far. Here is the question: It is quite possible that in some cases I lose everything wha…

---

## [Design Problems in Sink Connector](https://community.neo4j.com/t/design-problems-in-sink-connector/42910)

<div class="topic-metadata">

**Author:** [@pselvakumarmca](https://community.neo4j.com/u/pselvakumarmca)\
**Replies:** 1\
**Last updated:** [August 19, 2021, 8:38pm UTC](https://community.neo4j.com/t/design-problems-in-sink-connector/42910 "2021-08-19T20:38:31Z")

</div>

Im trying to use sink connectotr to sync data from mongodb. I receive relationship event as, "fullDocument": { "\_id": { "$oid": "6112376c9ade02082f500711" }, "\_relatedFromId": "6112376b9ade02082f5006f…

---

## [Connecting to a specific graph - neo4j spark connector](https://community.neo4j.com/t/connecting-to-a-specific-graph-neo4j-spark-connector/42416)

<div class="topic-metadata">

**Author:** [@poem\_daga](https://community.neo4j.com/u/poem_daga)\
**Replies:** 1\
**Last updated:** [August 9, 2021, 6:10am UTC](https://community.neo4j.com/t/connecting-to-a-specific-graph-neo4j-spark-connector/42416 "2021-08-09T06:10:00Z")

</div>

How to specify a specific graph to connect to with integration neo4j with spark? I have multiple graphs like: There is no option given in official document to specify which graph to use while reading/ writing data.

---

## [Neo4j sink instance not receiving events from Kafka (kafka Topic -\> Neo4j)](https://community.neo4j.com/t/neo4j-sink-instance-not-receiving-events-from-kafka-kafka-topic-neo4j/41381)

<div class="topic-metadata">

**Author:** [@bhuvana.rs](https://community.neo4j.com/u/bhuvana.rs)\
**Replies:** 2\
**Last updated:** [July 19, 2021, 6:08am UTC](https://community.neo4j.com/t/neo4j-sink-instance-not-receiving-events-from-kafka-kafka-topic-neo4j/41381 "2021-07-19T06:08:51Z")

</div>

Hello, I have configured neo4j and Kafka in K8s cluster (standalone). I'm using neo4j V4.1.3. I have configured neo4j as sink to consume JSON messages from Kafka topic. On sending the JSON message on a kafka topic, I d…

---

## [Not able to run 'Stream' Procedures](https://community.neo4j.com/t/not-able-to-run-stream-procedures/39776)

<div class="topic-metadata">

**Author:** [@vikash.kumar3](https://community.neo4j.com/u/vikash.kumar3)\
**Replies:** 1\
**Last updated:** [June 24, 2021, 12:14am UTC](https://community.neo4j.com/t/not-able-to-run-stream-procedures/39776 "2021-06-24T00:14:54Z")

</div>

Hi Everyone, I am not able to run the "Stream" procedures and getting this error: Failed to invoke procedure \`streams.consume\`: Caused by: java.lang.NoSuchMethodError: 'java.io.File org.neo4j.io.layout.DatabaseLayout.d…

---

## [How to call Apoc Procedures in Kafka Connect Cypher template?](https://community.neo4j.com/t/how-to-call-apoc-procedures-in-kafka-connect-cypher-template/37867)

<div class="topic-metadata">

**Author:** [@stefan.alschner](https://community.neo4j.com/u/stefan.alschner)\
**Replies:** 1\
**Last updated:** [May 13, 2021, 12:03pm UTC](https://community.neo4j.com/t/how-to-call-apoc-procedures-in-kafka-connect-cypher-template/37867 "2021-05-13T12:03:06Z")

</div>

Hi, I'm experimenting with with the Kafka / Confluent Neo4j-Sink-Connector. What i wanted to achieve is to send data from MongoDB via Kafka to Neo4j. So far I have succeeded in merging and creating simple nodes and rela…

---

## [How to act on changes to the database](https://community.neo4j.com/t/how-to-act-on-changes-to-the-database/35319)

<div class="topic-metadata">

**Author:** [@waters.simon](https://community.neo4j.com/u/waters.simon)\
**Replies:** 1\
**Last updated:** [March 16, 2021, 7:41am UTC](https://community.neo4j.com/t/how-to-act-on-changes-to-the-database/35319 "2021-03-16T07:41:44Z")

</div>

One design question for my app where I've not understood if there is a Neo4J way of achieving the desired effect is.... As web based users create or delete nodes or relationships, in Neo4j through a web app, I want to k…

---

## [Oracle RDBMS to neo4j graph incremental data transfer](https://community.neo4j.com/t/oracle-rdbms-to-neo4j-graph-incremental-data-transfer/30930)

<div class="topic-metadata">

**Author:** [@nar.kendre](https://community.neo4j.com/u/nar.kendre)\
**Replies:** 0\
**Last updated:** [December 24, 2020, 5:43pm UTC](https://community.neo4j.com/t/oracle-rdbms-to-neo4j-graph-incremental-data-transfer/30930 "2020-12-24T17:43:51Z")

</div>

Experts, We have done the infra setup to transfer the incremental data from Oracle 12c to Neo4j 3.5.5 (RDMS \<-\> Oracle Golden Gate \<-\> Big Data Gold Gate \<-\> Kafka \<-\> Neo4j stream \<-\> Neo4j) And able to see the updated…

---

## [Occur error when used to streams plugin in standalone mode](https://community.neo4j.com/t/occur-error-when-used-to-streams-plugin-in-standalone-mode/30852)

<div class="topic-metadata">

**Author:** [@admin7](https://community.neo4j.com/u/admin7)\
**Replies:** 0\
**Last updated:** [December 23, 2020, 8:42am UTC](https://community.neo4j.com/t/occur-error-when-used-to-streams-plugin-in-standalone-mode/30852 "2020-12-23T08:42:29Z")

</div>

Neo4j Version : neo4j-enterprise 4.2.1 neo4j-streams plugin version : 4.0.6 Server : AWS t3.medium x1 / ami-066e8505716e573cf (neo4j-enterprise-1-4.2.1-apoc 2020-11-30T04\_03\_27Z) I didn't use cloudformation. Neo4j d…

---

## [Neo4j Kafka Source configuration for initial load](https://community.neo4j.com/t/neo4j-kafka-source-configuration-for-initial-load/27937)

<div class="topic-metadata">

**Author:** [@basavarajdhanashetti](https://community.neo4j.com/u/basavarajdhanashetti)\
**Replies:** 1\
**Last updated:** [October 27, 2020, 11:29am UTC](https://community.neo4j.com/t/neo4j-kafka-source-configuration-for-initial-load/27937 "2020-10-27T11:29:13Z")

</div>

Hi, I was looking for configuration for the Neo4j in conf.ini file to load the existing data onto the Kafka topic. With configuration provided Kafka Connect Neo4j Connector User Guide - Neo4j Kafka Integration Docs , C…

---

## [Neo4j Streams consuming messages with GREAT LAG, and sudden bumps](https://community.neo4j.com/t/neo4j-streams-consuming-messages-with-great-lag-and-sudden-bumps/18552)

<div class="topic-metadata">

**Author:** [@elwosto](https://community.neo4j.com/u/elwosto)\
**Replies:** 1\
**Last updated:** [May 14, 2020, 12:35pm UTC](https://community.neo4j.com/t/neo4j-streams-consuming-messages-with-great-lag-and-sudden-bumps/18552 "2020-05-14T12:35:14Z")

</div>

Hi All, I have set up a Neo4j server with Kafka Sink/Source integration enabled. The Kafka servers is a standard Managed Kafka Service in AWS (3 brokers), and the Neo4j server is a docker container running in a EC2 Inst…

---

## [Is it possible to Retry processing a message a after race condition?](https://community.neo4j.com/t/is-it-possible-to-retry-processing-a-message-a-after-race-condition/18640)

<div class="topic-metadata">

**Author:** [@elwosto](https://community.neo4j.com/u/elwosto)\
**Replies:** 0\
**Last updated:** [May 14, 2020, 12:31pm UTC](https://community.neo4j.com/t/is-it-possible-to-retry-processing-a-message-a-after-race-condition/18640 "2020-05-14T12:31:05Z")

</div>

Hi All, TL;DR: I want to know if messages can be automatically reprocessed when errors/exceptions have been thrown during the topic handler's query I'm running Neo4j 3.5.14, with a Sink and Source enabled Kafka Integra…

[Next page](https://community.neo4j.com/c/integrations/stream-processing/22.md?page=1)
