Skip to content

Commit c87ea09

Browse files
committed
Fix PulsarIO
The main issue for current PulsarIO.read is * It is based on Pulsar Reader instead of PulsarConsumer, which then do not support acknowledgement * The while() block in reader DoFn would never return until topic termination, this basically means pipeline stuck * The restriction is on publishTime, and tryClaim assumes its ordering. This is not true. reader returning message is ordered on messageId. This is a wrong choice. Currently unresolved * PulsarMessage's coder implementation dropped message. This causes Data loss if the PulsarIO.read do not follow an immediate mapping * Tests are defunct and errors are suppressed, making them succeed spuriously Current PulsarIO.write is even more primitive. Pipeline expansion actually fails. It is not idempotent. Major fixes include * Allow Pulsar reader to have a timeout * Fix PulsarMessage and coder to include serializable fields from message * Fix mock client/reader and add a full read pipeline in test * Fix issues prevent PulsarIO.write from expanding. now it works minimally, that is publish every message received (at least once). * Working integration tests for read and write This has made PulsarIO.read minimally functionable. Although it won't split and can only run single thread. Going forward, we should re-implement reader DoFn based on Pulsar consumer. Thoughs rename the current DoFn to "NaiveReadFromPulsarDoFn"
1 parent 9651c64 commit c87ea09

17 files changed

Lines changed: 770 additions & 541 deletions

.github/workflows/beam_PreCommit_Java_Pulsar_IO_Direct.yml

Lines changed: 9 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -21,31 +21,13 @@ on:
2121
branches: ['master', 'release-*']
2222
paths:
2323
- "sdks/java/io/pulsar/**"
24-
- "sdks/java/io/common/**"
25-
- "sdks/java/core/src/main/**"
26-
- "build.gradle"
27-
- "buildSrc/**"
28-
- "gradle/**"
29-
- "gradle.properties"
30-
- "gradlew"
31-
- "gradle.bat"
32-
- "settings.gradle.kts"
3324
- ".github/workflows/beam_PreCommit_Java_Pulsar_IO_Direct.yml"
3425
pull_request_target:
3526
branches: ['master', 'release-*']
3627
paths:
3728
- "sdks/java/io/pulsar/**"
38-
- "sdks/java/io/common/**"
39-
- "sdks/java/core/src/main/**"
29+
- ".github/workflows/beam_PreCommit_Java_Pulsar_IO_Direct.yml"
4030
- 'release/trigger_all_tests.json'
41-
- '.github/trigger_files/beam_PreCommit_Java_Pulsar_IO_Direct.json'
42-
- "build.gradle"
43-
- "buildSrc/**"
44-
- "gradle/**"
45-
- "gradle.properties"
46-
- "gradlew"
47-
- "gradle.bat"
48-
- "settings.gradle.kts"
4931
issue_comment:
5032
types: [created]
5133
schedule:
@@ -110,6 +92,13 @@ jobs:
11092
arguments: |
11193
-PdisableSpotlessCheck=true \
11294
-PdisableCheckStyle=true \
95+
- name: run Pulsar IO IT script
96+
uses: ./.github/actions/gradle-command-self-hosted-action
97+
with:
98+
gradle-command: :sdks:java:io:pulsar:integrationTest
99+
arguments: |
100+
-PdisableSpotlessCheck=true \
101+
-PdisableCheckStyle=true \
113102
- name: Archive JUnit Test Results
114103
uses: actions/upload-artifact@v4
115104
if: ${{ !success() }}
@@ -135,4 +124,4 @@ jobs:
135124
if: always()
136125
with:
137126
name: Publish SpotBugs
138-
path: '**/build/reports/spotbugs/*.html'
127+
path: '**/build/reports/spotbugs/*.html'

build.gradle.kts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -355,6 +355,7 @@ tasks.register("javaioPreCommit") {
355355
dependsOn(":sdks:java:io:mqtt:build")
356356
dependsOn(":sdks:java:io:neo4j:build")
357357
dependsOn(":sdks:java:io:parquet:build")
358+
dependsOn(":sdks:java:io:pulsar:build")
358359
dependsOn(":sdks:java:io:rabbitmq:build")
359360
dependsOn(":sdks:java:io:redis:build")
360361
dependsOn(":sdks:java:io:rrio:build")

sdks/java/io/pulsar/build.gradle

Lines changed: 11 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -18,31 +18,32 @@
1818

1919
plugins { id 'org.apache.beam.module' }
2020
applyJavaNature(automaticModuleName: 'org.apache.beam.sdk.io.pulsar')
21+
enableJavaPerformanceTesting()
2122

2223
description = "Apache Beam :: SDKs :: Java :: IO :: Pulsar"
2324
ext.summary = "IO to read and write to Pulsar"
2425

25-
def pulsar_version = '2.8.2'
26+
def pulsar_version = '2.11.4'
2627

2728

2829
dependencies {
2930
implementation library.java.vendored_guava_32_1_2_jre
3031
implementation library.java.slf4j_api
3132
implementation library.java.joda_time
3233

33-
implementation "org.apache.pulsar:pulsar-client:$pulsar_version"
34-
implementation "org.apache.pulsar:pulsar-client-admin:$pulsar_version"
35-
permitUnusedDeclared "org.apache.pulsar:pulsar-client:$pulsar_version"
36-
permitUnusedDeclared "org.apache.pulsar:pulsar-client-admin:$pulsar_version"
37-
permitUsedUndeclared "org.apache.pulsar:pulsar-client-api:$pulsar_version"
38-
permitUsedUndeclared "org.apache.pulsar:pulsar-client-admin-api:$pulsar_version"
34+
implementation "org.apache.pulsar:pulsar-client-api:$pulsar_version"
35+
implementation "org.apache.pulsar:pulsar-client-admin-api:$pulsar_version"
36+
runtimeOnly "org.apache.pulsar:pulsar-client:$pulsar_version"
37+
runtimeOnly("org.apache.pulsar:pulsar-client-admin:$pulsar_version") {
38+
// To prevent a StackOverflow within Pulsar admin client because JUL -> SLF4J -> JUL
39+
exclude group: "org.slf4j", module: "jul-to-slf4j"
40+
}
3941

4042
implementation project(path: ":sdks:java:core", configuration: "shadow")
4143

42-
testImplementation library.java.jupiter_api
43-
testRuntimeOnly library.java.jupiter_engine
44+
testImplementation library.java.junit
45+
testRuntimeOnly library.java.slf4j_jdk14
4446
testRuntimeOnly project(path: ":runners:direct-java", configuration: "shadow")
4547
testImplementation "org.testcontainers:pulsar:1.15.3"
4648
testImplementation "org.assertj:assertj-core:2.9.1"
47-
4849
}

0 commit comments

Comments
 (0)