This is an automated email from the ASF dual-hosted git repository.
github-actions[bot] pushed a change to branch nightly-refs/heads/master
in repository https://gitbox.apache.org/repos/asf/beam.git
from c74609d71ee [Python] Support Lakehouse runtime catalog tables in
ReadFromBigQuery… (#40225)
add 0b6d68278b8 [KafkaIO] Use consumer position and lag to estimate end
offsets in KafkaUnboundedReader (#39830)
add ac516686b9f Bump github.com/proullon/ramsql from 0.1.4 to 0.1.5 in
/sdks (#40425)
add c5c59f1da69 Bump cloud.google.com/go/storage from 1.68.0 to 1.69.0 in
/sdks (#40426)
add 63a3479229b Update self-hosted runner versions to meet GitHub's
minimum runner version (#40403)
add d803b846618 [Python] Remove leftover "### Labels" debug log in the
BigQuery source (#40430)
add b456c68217d [Go SDK] Use deterministic protobuf marshaling in protox
helpers (#40385)
add b1b914b8bd0 Update Grafana version to 12.4.12 (#40434)
add 6fe144b3768 [BEAM-40264] Fix Spark MapState and SetState isEmpty after
removals (#40398)
add 32d54e821a6 [Spark] Run ValidatesRunner tests serially, runner metrics
are JVM wide (#40428)
add dfc4d667cd9 [Spark 4] Support stateful ParDo in Structured Streaming
via transformWithState (#40281)
add c164b38e924 Add a security model for Beam (#40421)
No new revisions were added by this update.
Summary of changes:
.../arc/variables.tf | 4 +-
.../self-hosted-linux/docker/Dockerfile | 2 +-
.../beam_PostCommit_Java_PVR_Spark3_Streaming.json | 2 +-
.../beam_PostCommit_Java_PVR_Spark4_Batch.json | 3 +-
.../beam_PostCommit_Java_PVR_Spark4_Streaming.json | 3 +-
...Commit_Java_PVR_Spark4_StructuredStreaming.json | 2 +-
.../beam_PostCommit_Java_PVR_Spark_Batch.json | 2 +-
...beam_PostCommit_Java_ValidatesRunner_Spark.json | 2 +-
...eam_PostCommit_Java_ValidatesRunner_Spark4.json | 2 +-
...a_ValidatesRunner_SparkStructuredStreaming.json | 4 +-
.test-infra/metrics/grafana/Dockerfile | 2 +-
.../beam/runners/core/StateInternalsTest.java | 10 +
.../translation/PipelineTranslatorStreaming.java | 65 +++-
.../StatefulParDoStreamingTranslator.java | 86 ++++++
.../streaming/state/BeamStatefulProcessor.java | 330 +++++++++++++++++++++
.../io/streaming/TestUnboundedSource.java | 16 +-
.../PipelineTranslatorStreamingTest.java | 97 +++++-
.../streaming/StatefulParDoStreamingTest.java | 242 +++++++++++++++
runners/spark/job-server/spark_job_server.gradle | 5 +
runners/spark/spark_runner.gradle | 5 +
.../spark/stateful/SparkStateInternals.java | 90 ++++--
.../spark/stateful/SparkTimerInternals.java | 5 +
.../translation/batch/DoFnRunnerFactory.java | 15 +-
.../batch/StatefulDoFnGroupFunction.java | 118 +-------
.../translation/batch/StatefulTaskRunner.java | 145 +++++++++
sdks/go.mod | 8 +-
sdks/go.sum | 20 +-
sdks/go/pkg/beam/core/util/protox/any.go | 4 +-
sdks/go/pkg/beam/core/util/protox/any_test.go | 40 +++
sdks/go/pkg/beam/core/util/protox/base64.go | 2 +-
sdks/go/pkg/beam/core/util/protox/protox.go | 2 +-
.../beam/sdk/io/kafka/KafkaUnboundedReader.java | 181 +++++------
.../org/apache/beam/sdk/io/kafka/KafkaIOTest.java | 27 --
.../org/apache/beam/sdk/io/kafka/KafkaMocks.java | 37 ---
sdks/python/apache_beam/io/gcp/bigquery.py | 1 -
website/www/site/content/en/security/_index.md | 93 ++++--
.../content/en/security/{_index.md => archive.md} | 13 +-
37 files changed, 1326 insertions(+), 359 deletions(-)
create mode 100644
runners/spark/4/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/StatefulParDoStreamingTranslator.java
create mode 100644
runners/spark/4/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/state/BeamStatefulProcessor.java
create mode 100644
runners/spark/4/src/test/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/StatefulParDoStreamingTest.java
create mode 100644
runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/StatefulTaskRunner.java
copy website/www/site/content/en/security/{_index.md => archive.md} (75%)