Skip to content

[BEAM-22] Support Unbounded PCollections in same-process execution - #3

Closed
tgroh wants to merge 3 commits into
apache:masterfrom
tgroh:in_process_pipeline_runner
Closed

[BEAM-22] Support Unbounded PCollections in same-process execution#3
tgroh wants to merge 3 commits into
apache:masterfrom
tgroh:in_process_pipeline_runner

Conversation

@tgroh

Copy link
Copy Markdown
Member

This completes the implementation of the InProcessPipelineRunner, which is capable of running over unbounded inputs.

Existing Unit Tests built on the DirectPipelineRunner all pass, with the exception of ParDoTest, which has a gap in immutability and encodability testing. This change does not migrate existing unit tests to run on the InProcessPipelineRunner.

@davorbonaci

Copy link
Copy Markdown
Member

Assignee: kennknowles
(we lack permission to set it on the pull request itself)

@tgroh
tgrohforce-pushed the in_process_pipeline_runner branch from 5feea5b to 5a15b31CompareMarch 1, 2016 19:28
tgroh added 3 commits March 1, 2016 11:52
This is the primary "global state" object for the evaluation of a
Pipeline using the InProcessPipelineRunner, and is responsible for
properly routing information about the state of the pipeline to
transform evaluators.
Remove the InProcessEvaluationContext from the InProcessPipelineRunner
class, and implement as a class directly. Fix associated imports.
This is responsible for scheduling transform evaluations and
communicating results back to the evaluation context. The executor
handle PTransforms that block arbitarily waiting for additional input.
Appropriately construct an evaluation context and executor, and start
the pipeline when run is called.
Implement runner-provided ExecutionContext and StepContext abstractions.
@davorbonaci

Copy link
Copy Markdown
Member

R: @kennknowles

@tgroh
tgrohforce-pushed the in_process_pipeline_runner branch from 5a15b31 to ccc368aCompareMarch 2, 2016 17:28
@tgroh

Copy link
Copy Markdown
MemberAuthor

Closing in favor of smaller requests

@tgrohtgroh closed this Mar 10, 2016
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Oct 14, 2016
cosmoskitten pushed a commit to cosmoskitten/beam that referenced this pull request Feb 21, 2017
pull latest code to check WatermarkTest
axelmagn referenced this pull request in axelmagn/beam Feb 15, 2018
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Mar 9, 2018
aljoscha referenced this pull request in aljoscha/beam Mar 12, 2018
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Apr 4, 2018
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Apr 5, 2018
mareksimunek referenced this pull request in mareksimunek/beam May 9, 2018
dmvk pushed a commit to dmvk/beam that referenced this pull request May 15, 2018
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Jun 4, 2018
…slation/BEAM-4410
[BEAM-4410] added BroadcastJoinTranslator
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Jun 7, 2018
* Update Schema.java
The type sets are final mutable types, but with immutable implementations. This doesn't communicate the desired semantic guarantees of the immutable implementation.
* Remove unneeded collection import.
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Oct 2, 2018
fixup! Read publish remote from command-line property
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Oct 2, 2018
swegner pushed a commit that referenced this pull request Oct 2, 2018
kennknowles pushed a commit that referenced this pull request Oct 16, 2018
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Oct 25, 2018
charlesccychen pushed a commit to cosmoskitten/beam that referenced this pull request Oct 29, 2018
VrishaliShah pushed a commit to VrishaliShah/beam that referenced this pull request Feb 5, 2020
update fork 3rd Dec 2019, 14:55PM
pabloem pushed a commit that referenced this pull request Feb 17, 2021
Debeziumio PoC (#7)
* New DebeziumIO class.
* Merge connector code
* DebeziumIO and MySqlConnector integrated.
* Added FormatFuntion param to Read builder on DebeziumIO.
* Added arguments checker to DebeziumIO.
* Add simple JSON mapper object (#1)
* Add simple JSON mapper object
* Fixed Mapper.
* Add SqlServer connector test
* Added PostgreSql Connector Test
PostgreSql now works with Json mapper
* Added PostgreSql Connector Test
PostgreSql now works with Json mapper
* Fixing MySQL schema DataException
Using file instead of schema should fix it
* MySQL Connector updated from 1.3.0 to 1.3.1
Co-authored-by: osvaldo-salinas <osvaldo.salinas@wizeline.com>
Co-authored-by: Carlos Dominguez <carlos.dominguez@carlos.dominguez>
Co-authored-by: Carlos Domínguez <carlos.dominguez@wizeline.com>
* Add debeziumio tests
* Debeziumio testing json mapper (#3)
* Some code refactors. Use a default DBHistory if not provided
* Add basic tests for Json mapper
* Debeziumio time restriction (#5)
* Add simple JSON mapper object
* Fixed Mapper.
* Add SqlServer connector test
* Added PostgreSql Connector Test
PostgreSql now works with Json mapper
* Added PostgreSql Connector Test
PostgreSql now works with Json mapper
* Fixing MySQL schema DataException
Using file instead of schema should fix it
* MySQL Connector updated from 1.3.0 to 1.3.1
* Some code refactors. Use a default DBHistory if not provided
* Adding based-time restriction
Stop polling after specified amount of time
* Add basic tests for Json mapper
* Adding new restriction
Uses a time-based restriction
* Adding optional restrcition
Uses an optional time-based restriction
Co-authored-by: juanitodread <juanitodread@gmail.com>
Co-authored-by: osvaldo-salinas <osvaldo.salinas@wizeline.com>
* Upgrade DebeziumIO connector (#4)
* Address comments (Change dependencies to testCompile, Set JsonMapper/Coder as default, refactors) (#8)
* Revert file
* Change dependencies to testCompile
* Move Counter sample to unit test
* Set JsonMapper as default mapper function
* Set String Coder as default coder when using JsonMapper
* Change logs from info to debug
* Debeziumio javadoc (#9)
* Adding javadoc
* Added some titles and examples
* Added SourceRecordJson doc
* Added Basic Connector doc
* Added KafkaSourceConsumer doc
* Javadoc cleanup
* Removing BasicConnector
No usages of this class were found overall
* Editing documentation
* Debeziumio fetched records restriction (#10)
* Adding javadoc
* Adding restriction by number of fetched records
Also adding a quick-fix for null value within SourceRecords
Minor fix on both MySQL and PostgreSQL Connectors Tests
* Run either by time or by number of records
* Added DebeziumOffsetTrackerTest
Tests both restrictions: By amount of time and by Number of records
* Removing comment
* DebeziumIO test for DB2. (#11)
* DebeziumIO test for DB2.
* DebeziumIO javadoc.
* Clean code:removed commented code lines on DebeziumIOConnectorTest.java
* Clean code:removing unused imports and using readAsJson().
Co-authored-by: Carlos Domínguez <74681048+carlosdominguezwl@users.noreply.github.com>
* Debezium limit records (now configurable) (#12)
* Adding javadoc
* Records Limit is now configurable
(It was fixed before)
* Debeziumio dockerize (#13)
* Add mysql docker container to tests
* Move debezium mysql integration test to its own file
* Add assertion to verify that the results contains a record.
* Debeziumio readme (#15)
* Adding javadoc
* Adding README file
* Add number of records configuration to the DebeziumIO component (#16)
* Code refactors (#17)
* Remove/ignore null warnings
* Remove DB2 code
* Remove docker dependency in DebeziumIO unit test and max number of recods to MySql integration test
* Change access modifiers accordingly
* Remove incomplete integration tests (Postgres and SqlServer)
* Add experimenal tag
* Debezium testing stoppable consumer (#18)
* Add try-catch-finally, stop SourceTask at finally.
* Fix warnings
* stopConsumer and processedRecords local variables removed. UT for task stop use case added
* Fix minor code style issue
Co-authored-by: juanitodread <juanitodread@gmail.com>
* Fix style issues (check, spotlessApply) (#19)
Co-authored-by: Osvaldo Salinas <osvaldo.salinas@osvaldo.salinas>
Co-authored-by: alejandro.maguey <alejandro.maguey@wizeline.com>
Co-authored-by: osvaldo-salinas <osvaldo.salinas@wizeline.com>
Co-authored-by: Carlos Dominguez <carlos.dominguez@carlos.dominguez>
Co-authored-by: Carlos Domínguez <carlos.dominguez@wizeline.com>
Co-authored-by: Carlos Domínguez <74681048+carlosdominguezwl@users.noreply.github.com>
Co-authored-by: Alejandro Maguey <alexmaguey1@gmail.com>
Co-authored-by: Hassan Reyes <hassanreyes@users.noreply.github.com>
Add missing apache license to README.md
Enabling integration test for DebeziumIO (#20)
Rename connector package cdc=>debezium. Update doc references (#21)
Fix code style on DebeziumIOMySqlConnectorIT
josiasrico added a commit to josiasrico/beam that referenced this pull request Mar 24, 2021
* Update Overview Process of Release Guide
Improved wording in first paragraphs.
Updated eleven steps to match with the sections. Added description to image of process.
* Update steps 4 and 7
Step 4: Added the rc variable in both (automatic and manual) steps and Attention notes. Step 7: Added a final step - Upload release candidate to PyPI.
usingh83 added a commit to usingh83/beam that referenced this pull request May 7, 2021
# This is the 1st commit message:
Java PreCommit failure fix
spotless failure fix
Java PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit refix
Java_Examples_Dataflow PreCommit fix
build failure corrected
Spotless check
Spotless check
reorganizing pipeline
delete the unused folder
Revert "Delete build.gradle"
This reverts commit c39a4e44
Delete build.gradle
don't need this file
adding comments and java docs, and removing unneeded dependencies.
Linting the project and making some stuff private
Reorganized and redefined to logic as per standard beam IO structure.
Lint the files.
Added changes for making the implementation more streamlined and understandable
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
# This is a combination of 15 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message apache#3:
Lint the files.
# This is the commit message apache#4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message apache#5:
Linting the project and making some stuff private
# This is the commit message apache#6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message apache#7:
delete the unused folder
# This is the commit message apache#8:
reorganizing pipeline
# This is the commit message apache#9:
Spotless check
# This is the commit message apache#10:
Spotless check
# This is the commit message apache#11:
build failure corrected
# This is the commit message apache#12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message apache#13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message apache#14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#15:
Java_Examples_Dataflow PreCommit assign nullable correctly
usingh83 added a commit to usingh83/beam that referenced this pull request May 13, 2021
# This is the 1st commit message:
# This is a combination of 2 commits.
# This is the 1st commit message:
Java PreCommit failure fix
spotless failure fix
Java PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit refix
Java_Examples_Dataflow PreCommit fix
build failure corrected
Spotless check
Spotless check
reorganizing pipeline
delete the unused folder
Revert "Delete build.gradle"
This reverts commit c39a4e44
Delete build.gradle
don't need this file
adding comments and java docs, and removing unneeded dependencies.
Linting the project and making some stuff private
Reorganized and redefined to logic as per standard beam IO structure.
Lint the files.
Added changes for making the implementation more streamlined and understandable
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
# This is a combination of 15 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message apache#3:
Lint the files.
# This is the commit message apache#4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message apache#5:
Linting the project and making some stuff private
# This is the commit message apache#6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message apache#7:
delete the unused folder
# This is the commit message apache#8:
reorganizing pipeline
# This is the commit message apache#9:
Spotless check
# This is the commit message apache#10:
Spotless check
# This is the commit message apache#11:
build failure corrected
# This is the commit message apache#12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message apache#13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message apache#14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#2:
# This is a combination of 3 commits.
# This is the 1st commit message:
Java PreCommit failure fix
spotless failure fix
Java PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit refix
Java_Examples_Dataflow PreCommit fix
build failure corrected
Spotless check
Spotless check
reorganizing pipeline
delete the unused folder
Revert "Delete build.gradle"
This reverts commit c39a4e44
Delete build.gradle
don't need this file
adding comments and java docs, and removing unneeded dependencies.
Linting the project and making some stuff private
Reorganized and redefined to logic as per standard beam IO structure.
Lint the files.
Added changes for making the implementation more streamlined and understandable
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
# This is a combination of 15 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message apache#3:
Lint the files.
# This is the commit message apache#4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message apache#5:
Linting the project and making some stuff private
# This is the commit message apache#6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message apache#7:
delete the unused folder
# This is the commit message apache#8:
reorganizing pipeline
# This is the commit message apache#9:
Spotless check
# This is the commit message apache#10:
Spotless check
# This is the commit message apache#11:
build failure corrected
# This is the commit message apache#12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message apache#13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message apache#14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#3:
# This is a combination of 16 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message apache#3:
Lint the files.
# This is the commit message apache#4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message apache#5:
Linting the project and making some stuff private
# This is the commit message apache#6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message apache#7:
delete the unused folder
# This is the commit message apache#8:
reorganizing pipeline
# This is the commit message apache#9:
Spotless check
# This is the commit message apache#10:
Spotless check
# This is the commit message apache#11:
build failure corrected
# This is the commit message apache#12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message apache#13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message apache#14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#16:
Java PreCommit assign nullable correctly
Java PreCommit assign nullable correctly
spotless failure fix
Java PreCommit failure fix
correcting the if checks
cleaning up and adding readme
spotless fixed
readme fixed and compileJava
fix
compileJava fix
compileJava fix now
spotless fix now
Java PreCommi fix
Java PreCommit fix
# This is a combination of 16 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message apache#3:
Lint the files.
# This is the commit message apache#4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message apache#5:
Linting the project and making some stuff private
# This is the commit message apache#6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message apache#7:
delete the unused folder
# This is the commit message apache#8:
reorganizing pipeline
# This is the commit message apache#9:
Spotless check
# This is the commit message apache#10:
Spotless check
# This is the commit message apache#11:
build failure corrected
# This is the commit message apache#12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message apache#13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message apache#14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#16:
Java PreCommit assign nullable correctly
Java PreCommit assign nullable correctly
spotless failure fix
Java PreCommit failure fix
correcting the if checks
cleaning up and adding readme
spotless fixed
readme fixed and compileJava
fix
compileJava fix
compileJava fix now
spotless fix now
Java PreCommi fix
Java PreCommit fix
# This is a combination of 3 commits.
# This is the 1st commit message:
Java PreCommit failure fix
spotless failure fix
Java PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit refix
Java_Examples_Dataflow PreCommit fix
build failure corrected
Spotless check
Spotless check
reorganizing pipeline
delete the unused folder
Revert "Delete build.gradle"
This reverts commit c39a4e44
Delete build.gradle
don't need this file
adding comments and java docs, and removing unneeded dependencies.
Linting the project and making some stuff private
Reorganized and redefined to logic as per standard beam IO structure.
Lint the files.
Added changes for making the implementation more streamlined and understandable
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
# This is a combination of 15 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message apache#3:
Lint the files.
# This is the commit message apache#4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message apache#5:
Linting the project and making some stuff private
# This is the commit message apache#6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message apache#7:
delete the unused folder
# This is the commit message apache#8:
reorganizing pipeline
# This is the commit message apache#9:
Spotless check
# This is the commit message apache#10:
Spotless check
# This is the commit message apache#11:
build failure corrected
# This is the commit message apache#12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message apache#13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message apache#14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#3:
# This is a combination of 16 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message apache#3:
Lint the files.
# This is the commit message apache#4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message apache#5:
Linting the project and making some stuff private
# This is the commit message apache#6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message apache#7:
delete the unused folder
# This is the commit message apache#8:
reorganizing pipeline
# This is the commit message apache#9:
Spotless check
# This is the commit message apache#10:
Spotless check
# This is the commit message apache#11:
build failure corrected
# This is the commit message apache#12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message apache#13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message apache#14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#16:
Java PreCommit assign nullable correctly
Java PreCommit assign nullable correctly
spotless failure fix
Java PreCommit failure fix
correcting the if checks
cleaning up and adding readme
spotless fixed
readme fixed and compileJava
fix
compileJava fix
compileJava fix now
spotless fix now
Java PreCommi fix
Java PreCommit fix
# This is a combination of 16 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message apache#2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message apache#3:
Lint the files.
# This is the commit message apache#4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message apache#5:
Linting the project and making some stuff private
# This is the commit message apache#6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message apache#7:
delete the unused folder
# This is the commit message apache#8:
reorganizing pipeline
# This is the commit message apache#9:
Spotless check
# This is the commit message apache#10:
Spotless check
# This is the commit message apache#11:
build failure corrected
# This is the commit message apache#12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message apache#13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message apache#14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message apache#16:
Java PreCommit assign nullable correctly
Java PreCommit assign nullable correctly
spotless failure fix
Java PreCommit failure fix
correcting the if checks
cleaning up and adding readme
spotless fixed
readme fixed and compileJava
fix
compileJava fix
compileJava fix now
spotless fix now
Java PreCommi fix
Java PreCommit fix
Final Commit with all changes
Added unit test
adding examples for usage
usage for TwitterIO added and Java PreCommit failure fix
Spotless PreCommit failure fix
pabloem pushed a commit that referenced this pull request May 18, 2021
…eams data from twitter
* # This is a combination of 2 commits.
# This is the 1st commit message:
Java PreCommit failure fix
spotless failure fix
Java PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit refix
Java_Examples_Dataflow PreCommit fix
build failure corrected
Spotless check
Spotless check
reorganizing pipeline
delete the unused folder
Revert "Delete build.gradle"
This reverts commit c39a4e44
Delete build.gradle
don't need this file
adding comments and java docs, and removing unneeded dependencies.
Linting the project and making some stuff private
Reorganized and redefined to logic as per standard beam IO structure.
Lint the files.
Added changes for making the implementation more streamlined and understandable
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
# This is a combination of 15 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message #3:
Lint the files.
# This is the commit message #4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message #5:
Linting the project and making some stuff private
# This is the commit message #6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message #7:
delete the unused folder
# This is the commit message #8:
reorganizing pipeline
# This is the commit message #9:
Spotless check
# This is the commit message #10:
Spotless check
# This is the commit message #11:
build failure corrected
# This is the commit message #12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message #13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message #14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #15:
Java_Examples_Dataflow PreCommit assign nullable correctly
* # This is a combination of 2 commits.
# This is the 1st commit message:
# This is a combination of 2 commits.
# This is the 1st commit message:
Java PreCommit failure fix
spotless failure fix
Java PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit refix
Java_Examples_Dataflow PreCommit fix
build failure corrected
Spotless check
Spotless check
reorganizing pipeline
delete the unused folder
Revert "Delete build.gradle"
This reverts commit c39a4e44
Delete build.gradle
don't need this file
adding comments and java docs, and removing unneeded dependencies.
Linting the project and making some stuff private
Reorganized and redefined to logic as per standard beam IO structure.
Lint the files.
Added changes for making the implementation more streamlined and understandable
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
# This is a combination of 15 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message #3:
Lint the files.
# This is the commit message #4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message #5:
Linting the project and making some stuff private
# This is the commit message #6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message #7:
delete the unused folder
# This is the commit message #8:
reorganizing pipeline
# This is the commit message #9:
Spotless check
# This is the commit message #10:
Spotless check
# This is the commit message #11:
build failure corrected
# This is the commit message #12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message #13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message #14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #2:
# This is a combination of 3 commits.
# This is the 1st commit message:
Java PreCommit failure fix
spotless failure fix
Java PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit refix
Java_Examples_Dataflow PreCommit fix
build failure corrected
Spotless check
Spotless check
reorganizing pipeline
delete the unused folder
Revert "Delete build.gradle"
This reverts commit c39a4e44
Delete build.gradle
don't need this file
adding comments and java docs, and removing unneeded dependencies.
Linting the project and making some stuff private
Reorganized and redefined to logic as per standard beam IO structure.
Lint the files.
Added changes for making the implementation more streamlined and understandable
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
# This is a combination of 15 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message #3:
Lint the files.
# This is the commit message #4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message #5:
Linting the project and making some stuff private
# This is the commit message #6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message #7:
delete the unused folder
# This is the commit message #8:
reorganizing pipeline
# This is the commit message #9:
Spotless check
# This is the commit message #10:
Spotless check
# This is the commit message #11:
build failure corrected
# This is the commit message #12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message #13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message #14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #3:
# This is a combination of 16 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message #3:
Lint the files.
# This is the commit message #4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message #5:
Linting the project and making some stuff private
# This is the commit message #6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message #7:
delete the unused folder
# This is the commit message #8:
reorganizing pipeline
# This is the commit message #9:
Spotless check
# This is the commit message #10:
Spotless check
# This is the commit message #11:
build failure corrected
# This is the commit message #12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message #13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message #14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #16:
Java PreCommit assign nullable correctly
Java PreCommit assign nullable correctly
spotless failure fix
Java PreCommit failure fix
correcting the if checks
cleaning up and adding readme
spotless fixed
readme fixed and compileJava
fix
compileJava fix
compileJava fix now
spotless fix now
Java PreCommi fix
Java PreCommit fix
# This is a combination of 16 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message #3:
Lint the files.
# This is the commit message #4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message #5:
Linting the project and making some stuff private
# This is the commit message #6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message #7:
delete the unused folder
# This is the commit message #8:
reorganizing pipeline
# This is the commit message #9:
Spotless check
# This is the commit message #10:
Spotless check
# This is the commit message #11:
build failure corrected
# This is the commit message #12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message #13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message #14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #16:
Java PreCommit assign nullable correctly
Java PreCommit assign nullable correctly
spotless failure fix
Java PreCommit failure fix
correcting the if checks
cleaning up and adding readme
spotless fixed
readme fixed and compileJava
fix
compileJava fix
compileJava fix now
spotless fix now
Java PreCommi fix
Java PreCommit fix
# This is a combination of 3 commits.
# This is the 1st commit message:
Java PreCommit failure fix
spotless failure fix
Java PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit assign nullable correctly
Java_Examples_Dataflow PreCommit refix
Java_Examples_Dataflow PreCommit fix
build failure corrected
Spotless check
Spotless check
reorganizing pipeline
delete the unused folder
Revert "Delete build.gradle"
This reverts commit c39a4e44
Delete build.gradle
don't need this file
adding comments and java docs, and removing unneeded dependencies.
Linting the project and making some stuff private
Reorganized and redefined to logic as per standard beam IO structure.
Lint the files.
Added changes for making the implementation more streamlined and understandable
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
# This is a combination of 15 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message #3:
Lint the files.
# This is the commit message #4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message #5:
Linting the project and making some stuff private
# This is the commit message #6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message #7:
delete the unused folder
# This is the commit message #8:
reorganizing pipeline
# This is the commit message #9:
Spotless check
# This is the commit message #10:
Spotless check
# This is the commit message #11:
build failure corrected
# This is the commit message #12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message #13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message #14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #3:
# This is a combination of 16 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message #3:
Lint the files.
# This is the commit message #4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message #5:
Linting the project and making some stuff private
# This is the commit message #6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message #7:
delete the unused folder
# This is the commit message #8:
reorganizing pipeline
# This is the commit message #9:
Spotless check
# This is the commit message #10:
Spotless check
# This is the commit message #11:
build failure corrected
# This is the commit message #12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message #13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message #14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #16:
Java PreCommit assign nullable correctly
Java PreCommit assign nullable correctly
spotless failure fix
Java PreCommit failure fix
correcting the if checks
cleaning up and adding readme
spotless fixed
readme fixed and compileJava
fix
compileJava fix
compileJava fix now
spotless fix now
Java PreCommi fix
Java PreCommit fix
# This is a combination of 16 commits.
# This is the 1st commit message:
Added a connector that streams data from twitter using a Standard Twitter app.
# This is the commit message #2:
Added changes for making the implementation more streamlined and understandable
# This is the commit message #3:
Lint the files.
# This is the commit message #4:
Reorganized and redefined to logic as per standard beam IO structure.
# This is the commit message #5:
Linting the project and making some stuff private
# This is the commit message #6:
adding comments and java docs, and removing unneeded dependencies.
# This is the commit message #7:
delete the unused folder
# This is the commit message #8:
reorganizing pipeline
# This is the commit message #9:
Spotless check
# This is the commit message #10:
Spotless check
# This is the commit message #11:
build failure corrected
# This is the commit message #12:
Java_Examples_Dataflow PreCommit fix
# This is the commit message #13:
Java_Examples_Dataflow PreCommit refix
# This is the commit message #14:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #15:
Java_Examples_Dataflow PreCommit assign nullable correctly
# This is the commit message #16:
Java PreCommit assign nullable correctly
Java PreCommit assign nullable correctly
spotless failure fix
Java PreCommit failure fix
correcting the if checks
cleaning up and adding readme
spotless fixed
readme fixed and compileJava
fix
compileJava fix
compileJava fix now
spotless fix now
Java PreCommi fix
Java PreCommit fix
Final Commit with all changes
Added unit test
adding examples for usage
usage for TwitterIO added and Java PreCommit failure fix
Spotless PreCommit failure fix
* Unit test for multiple config added, and beautification
* Spotless apply fixed
* Removing redundant comments
* Removing newly added test
* adding newly added test back
PawasChhokra pushed a commit to PawasChhokra/beam that referenced this pull request Aug 27, 2021
* Create linkedin-publish.yml
Testing publishing with github actions
* Fix li_maven_publish workflow
* Update version numbers
* Update python versions
PawasChhokra pushed a commit to PawasChhokra/beam that referenced this pull request Sep 15, 2021
* Create linkedin-publish.yml
Testing publishing with github actions
* Fix li_maven_publish workflow
* Update version numbers
* Update python versions
pabloem pushed a commit that referenced this pull request Jan 21, 2022
Adds ReadChangeStreamPartitionDoFn, which is an SDF to read partitions
from change streams and process them accordingly. This component
receives a change stream name, a partition, a start time and an end time
to query. It then initiates a change stream query with the received
parameters.
Within a change stream, 3 types of records can be received:
1. A Data record
2. A Heartbeat record
3. A Child partitions record
Upon receiving #1, the function updates the watermark with the record's
commit timestamp and emits the record into the output PCollection.
Upon receiving #2, the function updates the watermark with the record's
timestamp, but it does not emit any record into the PCollection.
Upon receiving #3, the function updates the watermark with the record's
timestamp and writes the new child partitions into the metadata table.
These partitions will be later scheduled by the DetectNewPartitions
component.
Once the change stream query for the element partition finishes, it
marks the partition as finished in the metadata table and terminates.
hengfengli referenced this pull request in hengfengli/beam Mar 21, 2022
* feat: Pipeline Initialization
* Use SpannerAccessor
robertwb pushed a commit that referenced this pull request May 6, 2022
add build and clean script to compile ts
sjvanrossum pushed a commit to sjvanrossum/beam that referenced this pull request May 17, 2023
- Fix GroupByKey test by using KV instead of tuple.
- Fix data race when writing to SERIALIZED_FNS.
- Enable ensure_assert_fails test which was presumably previously causing
spurious failures due to data race.
sjvanrossum pushed a commit to sjvanrossum/beam that referenced this pull request May 17, 2023
@rajkguptrajkgupt mentioned this pull request Apr 22, 2024
3 tasks
shahar1 pushed a commit to shahar1/beam that referenced this pull request Jan 18, 2025
- Fix GroupByKey test by using KV instead of tuple.
- Fix data race when writing to SERIALIZED_FNS.
- Enable ensure_assert_fails test which was presumably previously causing
spurious failures due to data race.
Abacn added a commit that referenced this pull request Jan 23, 2025
Abacn added a commit that referenced this pull request Jan 23, 2025
* Revert "[Dataflow Streaming Appliance] Fix per key commit size validation (#3…"
This reverts commit edc4766.
* keep main src change
VardhanThigle pushed a commit to VardhanThigle/beam that referenced this pull request Mar 21, 2025
* Revert "[Dataflow Streaming Appliance] Fix per key commit size validation (apache#3…"
This reverts commit edc4766.
* keep main src change
shunping pushed a commit that referenced this pull request May 30, 2026
ash6898 pushed a commit to ash6898/beam that referenced this pull request Jun 29, 2026
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@tgroh@davorbonaci