This is an automated email from the ASF dual-hosted git repository.

shunping pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 7b9380b1c25 Buffer BufferedLogger by newline to avoid log splitting 
(#39288)
7b9380b1c25 is described below

commit 7b9380b1c2520766be6b3adea3fd153fb50a5427
Author: Shunping Huang <[email protected]>
AuthorDate: Thu Jul 30 17:35:39 2026 -0400

    Buffer BufferedLogger by newline to avoid log splitting (#39288)
    
    * Buffer BufferedLogger by newline to avoid log splitting
    
    * Update BufferedLogger.Write to search for newlines and accumulate
      partial log lines in the builder instead of immediately emitting
      them as separate entries.
    * Update FlushAtDebug and FlushAtError to flush any remaining
      trailing text in the builder on exit or flush events.
    * Fix bug in buffered_logging_test.go where log list assertions
      only verified the first element of logCatcher.msgs instead of
      checking all gathered log messages.
    * Add TestBufferedLogger/partial_write_splitting to verify correct
      chunked write buffering and line assembly behavior.
    
    * Call BuferredLogger.Printf so the deps will not be separate lines.
    
    * Address comments
    
    * Refactor BufferedLogger flushing and container command execution
    
    - Add Flush(ctx, err) to BufferedLogger to automatically flush at
      ERROR on failure or DEBUG on success.
    - Add executeWithLogger and executeWithOutput helpers to eliminate
      repetitive flush boilerplate.
    - Unify runtime dependency log outputs into single log entries.
    - Add unit tests and docstrings for BufferedLogger.
    
    * Address comments.
---
 sdks/go/container/tools/buffered_logging.go      |  64 +++++++--
 sdks/go/container/tools/buffered_logging_test.go | 168 +++++++++++++++++++++--
 sdks/python/container/boot.go                    |  23 +---
 sdks/python/container/piputil.go                 |  39 +++---
 4 files changed, 240 insertions(+), 54 deletions(-)

diff --git a/sdks/go/container/tools/buffered_logging.go 
b/sdks/go/container/tools/buffered_logging.go
index a0937b8eb14..333b136e16f 100644
--- a/sdks/go/container/tools/buffered_logging.go
+++ b/sdks/go/container/tools/buffered_logging.go
@@ -16,6 +16,7 @@
 package tools
 
 import (
+       "bytes"
        "context"
        "log"
        "os"
@@ -52,23 +53,59 @@ func NewBufferedLoggerWithFlushInterval(ctx 
context.Context, logger *Logger, int
        return &BufferedLogger{logger: logger, lastFlush: time.Now(), 
flushInterval: interval, periodicFlushContext: ctx, now: time.Now}
 }
 
-// Write implements the io.Writer interface, converting input to a string
-// and storing it in the BufferedLogger's buffer. If a logger is not provided,
-// the output is sent directly to os.Stderr.
+// Write implements the io.Writer interface. It buffers byte streams 
line-by-line
+// into memory and flushes periodically or upon calling Flush(), 
FlushAtError(), or
+// FlushAtDebug(). It is used primarily to redirect stdout/stderr of 
subprocesses or
+// standard Go log output. If a logger is not provided, the output is sent 
directly to os.Stderr.
 func (b *BufferedLogger) Write(p []byte) (int, error) {
        if b.logger == nil {
                return os.Stderr.Write(p)
        }
-       n, err := b.builder.Write(p)
+
        if b.logs == nil {
                b.logs = make([]string, 0, initialLogSize)
        }
-       b.logs = append(b.logs, b.builder.String())
-       b.builder.Reset()
+
+       start := 0
+       for {
+               // Look for the next newline in the incoming byte slice directly
+               nl := bytes.IndexByte(p[start:], '\n')
+               if nl == -1 {
+                       break
+               }
+
+               // Write the segment up to the newline into the builder
+               b.builder.Write(p[start : start+nl])
+
+               // The builder now contains any previous partial line + the 
current complete segment
+               b.logs = append(b.logs, strings.TrimSuffix(b.builder.String(), 
"\r"))
+               b.builder.Reset()
+
+               start += nl + 1
+       }
+
+       // Buffer any remaining bytes that didn't end in a newline
+       if start < len(p) {
+               b.builder.Write(p[start:])
+       }
+
        if b.now().Sub(b.lastFlush) > b.flushInterval {
                b.FlushAtDebug(b.periodicFlushContext)
        }
-       return n, err
+
+       return len(p), nil
+}
+
+// Flush flushes the contents of the buffer to the logging service.
+// If err is non-nil, it flushes at Error severity; otherwise it flushes at 
Debug severity.
+// It returns the provided error.
+func (b *BufferedLogger) Flush(ctx context.Context, err error) error {
+       if err != nil {
+               b.FlushAtError(ctx)
+       } else {
+               b.FlushAtDebug(ctx)
+       }
+       return err
 }
 
 // FlushAtError flushes the contents of the buffer to the logging
@@ -77,6 +114,10 @@ func (b *BufferedLogger) FlushAtError(ctx context.Context) {
        if b.logger == nil {
                return
        }
+       if b.builder.Len() > 0 {
+               b.logs = append(b.logs, strings.TrimSuffix(b.builder.String(), 
"\r"))
+               b.builder.Reset()
+       }
        for _, message := range b.logs {
                b.logger.Errorf(ctx, "%s", message)
        }
@@ -90,6 +131,10 @@ func (b *BufferedLogger) FlushAtDebug(ctx context.Context) {
        if b.logger == nil {
                return
        }
+       if b.builder.Len() > 0 {
+               b.logs = append(b.logs, strings.TrimSuffix(b.builder.String(), 
"\r"))
+               b.builder.Reset()
+       }
        for _, message := range b.logs {
                b.logger.Printf(ctx, "%s", message)
        }
@@ -97,8 +142,9 @@ func (b *BufferedLogger) FlushAtDebug(ctx context.Context) {
        b.lastFlush = time.Now()
 }
 
-// Prints directly to the logging service. If the logger is nil, prints 
directly to the
-// console. Used for the container pre-build workflow.
+// Printf directly writes formatted messages to the underlying logger/service,
+// bypassing line buffering. If the logger is nil, it prints directly to the
+// console. Used for direct informational logs and the container pre-build 
workflow.
 func (b *BufferedLogger) Printf(ctx context.Context, format string, args 
...any) {
        if b.logger == nil {
                log.Printf(format, args...)
diff --git a/sdks/go/container/tools/buffered_logging_test.go 
b/sdks/go/container/tools/buffered_logging_test.go
index 9f542d2d5ab..314d496475a 100644
--- a/sdks/go/container/tools/buffered_logging_test.go
+++ b/sdks/go/container/tools/buffered_logging_test.go
@@ -17,21 +17,48 @@ package tools
 
 import (
        "context"
+       "errors"
        "testing"
        "time"
 
        fnpb "github.com/apache/beam/sdks/v2/go/pkg/beam/model/fnexecution_v1"
 )
 
+func getAllLogEntries(catcher *logCatcher) []*fnpb.LogEntry {
+       var entries []*fnpb.LogEntry
+       for _, list := range catcher.msgs {
+               entries = append(entries, list.GetLogEntries()...)
+       }
+       return entries
+}
+
 func TestBufferedLogger(t *testing.T) {
        ctx := context.Background()
 
+       t.Run("printf", func(t *testing.T) {
+               catcher := &logCatcher{}
+               l := &Logger{client: catcher}
+               bl := NewBufferedLogger(l)
+
+               bl.Printf(ctx, "test message")
+
+               received := catcher.msgs[0].GetLogEntries()[0]
+
+               if got, want := received.Message, "test message"; got != want {
+                       t.Errorf("got message %q, want %q", got, want)
+               }
+
+               if got, want := received.Severity, 
fnpb.LogEntry_Severity_DEBUG; got != want {
+                       t.Errorf("got severity %v, want %v", got, want)
+               }
+       })
+
        t.Run("write", func(t *testing.T) {
                catcher := &logCatcher{}
                l := &Logger{client: catcher}
                bl := NewBufferedLogger(l)
 
-               message := []byte("test message")
+               message := []byte("test message\n")
                n, err := bl.Write(message)
                if err != nil {
                        t.Errorf("got error %v", err)
@@ -77,7 +104,8 @@ func TestBufferedLogger(t *testing.T) {
                l := &Logger{client: catcher}
                bl := NewBufferedLogger(l)
 
-               messages := []string{"foo", "bar", "baz"}
+               messages := []string{"foo\n", "bar\n", "baz\n"}
+               expected := []string{"foo", "bar", "baz"}
 
                for _, message := range messages {
                        messBytes := []byte(message)
@@ -93,10 +121,14 @@ func TestBufferedLogger(t *testing.T) {
 
                bl.FlushAtDebug(ctx)
 
-               received := catcher.msgs[0].GetLogEntries()
+               received := getAllLogEntries(catcher)
+
+               if got, want := len(received), len(expected); got != want {
+                       t.Fatalf("expected %d log entries received, got %d", 
want, got)
+               }
 
                for i, message := range received {
-                       if got, want := message.Message, messages[i]; got != 
want {
+                       if got, want := message.Message, expected[i]; got != 
want {
                                t.Errorf("got message %q, want %q", got, want)
                        }
 
@@ -139,7 +171,8 @@ func TestBufferedLogger(t *testing.T) {
                l := &Logger{client: catcher}
                bl := NewBufferedLogger(l)
 
-               messages := []string{"foo", "bar", "baz"}
+               messages := []string{"foo\n", "bar\n", "baz\n"}
+               expected := []string{"foo", "bar", "baz"}
 
                for _, message := range messages {
                        messBytes := []byte(message)
@@ -155,10 +188,14 @@ func TestBufferedLogger(t *testing.T) {
 
                bl.FlushAtError(ctx)
 
-               received := catcher.msgs[0].GetLogEntries()
+               received := getAllLogEntries(catcher)
+
+               if got, want := len(received), len(expected); got != want {
+                       t.Fatalf("expected %d log entries received, got %d", 
want, got)
+               }
 
                for i, message := range received {
-                       if got, want := message.Message, messages[i]; got != 
want {
+                       if got, want := message.Message, expected[i]; got != 
want {
                                t.Errorf("got message %q, want %q", got, want)
                        }
 
@@ -168,6 +205,55 @@ func TestBufferedLogger(t *testing.T) {
                }
        })
 
+       t.Run("flush with nil error", func(t *testing.T) {
+               catcher := &logCatcher{}
+               l := &Logger{client: catcher}
+               bl := NewBufferedLogger(l)
+
+               message := []byte("success message\n")
+               _, err := bl.Write(message)
+               if err != nil {
+                       t.Fatalf("unexpected write error: %v", err)
+               }
+
+               if gotErr := bl.Flush(ctx, nil); gotErr != nil {
+                       t.Errorf("Flush(ctx, nil) returned error %v, want nil", 
gotErr)
+               }
+
+               received := catcher.msgs[0].GetLogEntries()[0]
+               if got, want := received.Message, "success message"; got != 
want {
+                       t.Errorf("got message %q, want %q", got, want)
+               }
+               if got, want := received.Severity, 
fnpb.LogEntry_Severity_DEBUG; got != want {
+                       t.Errorf("got severity %v, want %v", got, want)
+               }
+       })
+
+       t.Run("flush with non-nil error", func(t *testing.T) {
+               catcher := &logCatcher{}
+               l := &Logger{client: catcher}
+               bl := NewBufferedLogger(l)
+
+               message := []byte("error message\n")
+               _, err := bl.Write(message)
+               if err != nil {
+                       t.Fatalf("unexpected write error: %v", err)
+               }
+
+               originalErr := errors.New("command failed")
+               if gotErr := bl.Flush(ctx, originalErr); gotErr != originalErr {
+                       t.Errorf("Flush(ctx, err) returned %v, want %v", 
gotErr, originalErr)
+               }
+
+               received := catcher.msgs[0].GetLogEntries()[0]
+               if got, want := received.Message, "error message"; got != want {
+                       t.Errorf("got message %q, want %q", got, want)
+               }
+               if got, want := received.Severity, 
fnpb.LogEntry_Severity_ERROR; got != want {
+                       t.Errorf("got severity %v, want %v", got, want)
+               }
+       })
+
        t.Run("direct print", func(t *testing.T) {
                catcher := &logCatcher{}
                l := &Logger{client: catcher}
@@ -195,7 +281,8 @@ func TestBufferedLogger(t *testing.T) {
                startTime := time.Now()
                bl.now = func() time.Time { return startTime }
 
-               messages := []string{"foo", "bar"}
+               messages := []string{"foo\n", "bar\n"}
+               expected := []string{"foo", "bar"}
 
                for i, message := range messages {
                        if i > 1 {
@@ -212,7 +299,8 @@ func TestBufferedLogger(t *testing.T) {
                        }
                }
 
-               lastMessage := "baz"
+               lastMessage := "baz\n"
+               expected = append(expected, "baz")
                bl.now = func() time.Time { return startTime.Add(6 * 
time.Second) }
                messBytes := []byte(lastMessage)
                n, err := bl.Write(messBytes)
@@ -225,11 +313,14 @@ func TestBufferedLogger(t *testing.T) {
                }
 
                // Type should have auto-flushed at debug after the third 
message
-               received := catcher.msgs[0].GetLogEntries()
-               messages = append(messages, lastMessage)
+               received := getAllLogEntries(catcher)
+
+               if got, want := len(received), len(expected); got != want {
+                       t.Fatalf("expected %d log entries received, got %d", 
want, got)
+               }
 
                for i, message := range received {
-                       if got, want := message.Message, messages[i]; got != 
want {
+                       if got, want := message.Message, expected[i]; got != 
want {
                                t.Errorf("got message %q, want %q", got, want)
                        }
 
@@ -238,4 +329,57 @@ func TestBufferedLogger(t *testing.T) {
                        }
                }
        })
+
+       t.Run("partial write splitting", func(t *testing.T) {
+               catcher := &logCatcher{}
+               l := &Logger{client: catcher}
+               bl := NewBufferedLogger(l)
+
+               // Write a partial line
+               n, err := bl.Write([]byte("hello "))
+               if err != nil {
+                       t.Errorf("got error %v", err)
+               }
+               if n != 6 {
+                       t.Errorf("got %d, want 6", n)
+               }
+               if len(bl.logs) != 0 {
+                       t.Errorf("expected no logs buffered yet, got %d", 
len(bl.logs))
+               }
+
+               // Write remainder and a second line
+               n, err = bl.Write([]byte("world\nline2\npartial"))
+               if err != nil {
+                       t.Errorf("got error %v", err)
+               }
+               if n != 19 {
+                       t.Errorf("got %d, want 19", n)
+               }
+
+               if got, want := len(bl.logs), 2; got != want {
+                       t.Errorf("expected 2 logs buffered, got %d", got)
+               }
+               if got, want := bl.logs[0], "hello world"; got != want {
+                       t.Errorf("got %q, want %q", got, want)
+               }
+               if got, want := bl.logs[1], "line2"; got != want {
+                       t.Errorf("got %q, want %q", got, want)
+               }
+
+               // Flush should flush the final partial message
+               bl.FlushAtDebug(ctx)
+               received := getAllLogEntries(catcher)
+               if got, want := len(received), 3; got != want {
+                       t.Fatalf("expected 3 log entries received, got %d", got)
+               }
+               if got, want := received[0].Message, "hello world"; got != want 
{
+                       t.Errorf("got message %q, want %q", got, want)
+               }
+               if got, want := received[1].Message, "line2"; got != want {
+                       t.Errorf("got message %q, want %q", got, want)
+               }
+               if got, want := received[2].Message, "partial"; got != want {
+                       t.Errorf("got message %q, want %q", got, want)
+               }
+       })
 }
diff --git a/sdks/python/container/boot.go b/sdks/python/container/boot.go
index 958fd46904a..364a614b4e8 100644
--- a/sdks/python/container/boot.go
+++ b/sdks/python/container/boot.go
@@ -571,19 +571,13 @@ func logRuntimeDependencies(ctx context.Context, 
bufLogger *tools.BufferedLogger
        if err != nil {
                return err
        }
-       bufLogger.Printf(ctx, "Python version in %s:", phase)
-       args := []string{"--version"}
-       if err := execx.ExecuteEnvWithIO(nil, os.Stdin, bufLogger, bufLogger, 
pythonVersion, args...); err != nil {
-               bufLogger.FlushAtError(ctx)
-       } else {
-               bufLogger.FlushAtDebug(ctx)
+       if out, err := executeWithOutput(ctx, bufLogger, pythonVersion, 
"--version"); err == nil {
+               bufLogger.Printf(ctx, "Python version in %s: %s", phase, 
strings.TrimSpace(string(out)))
        }
-       bufLogger.Printf(ctx, "Dependencies in %s:", phase)
-       args = []string{"-m", "pip", "freeze", "--all"}
-       if err := execx.ExecuteEnvWithIO(nil, os.Stdin, bufLogger, bufLogger, 
pythonVersion, args...); err != nil {
-               bufLogger.FlushAtError(ctx)
-       } else {
-               bufLogger.FlushAtDebug(ctx)
+
+       args := []string{"-m", "pip", "freeze", "--all"}
+       if out, err := executeWithOutput(ctx, bufLogger, pythonVersion, 
args...); err == nil {
+               bufLogger.Printf(ctx, "Dependencies in %s:\n%s", phase, 
string(out))
        }
        return nil
 }
@@ -591,7 +585,6 @@ func logRuntimeDependencies(ctx context.Context, bufLogger 
*tools.BufferedLogger
 // logSubmissionEnvDependencies logs the python dependencies
 // installed in the submission environment.
 func logSubmissionEnvDependencies(ctx context.Context, bufLogger 
*tools.BufferedLogger, dir string) error {
-       bufLogger.Printf(ctx, "Dependencies in submission environment:")
        // path for submission environment dependencies should match with the
        // one defined in apache_beam/runners/portability/stager.py.
        filename := filepath.Join(dir, 
"submission_environment_dependencies.txt")
@@ -599,8 +592,6 @@ func logSubmissionEnvDependencies(ctx context.Context, 
bufLogger *tools.Buffered
        if err != nil {
                return err
        }
-       bufLogger.Printf(ctx, "%s", string(content))
+       bufLogger.Printf(ctx, "Dependencies in submission environment:\n%s", 
string(content))
        return nil
 }
-
-
diff --git a/sdks/python/container/piputil.go b/sdks/python/container/piputil.go
index 2024c16dde5..5cef517fc42 100644
--- a/sdks/python/container/piputil.go
+++ b/sdks/python/container/piputil.go
@@ -41,6 +41,24 @@ var (
 const pipLogFlushInterval time.Duration = 15 * time.Second
 const unrecoverableURL string = 
"https://beam.apache.org/documentation/sdks/python-unrecoverable-errors/index.html#pip-dependency-resolution-failures";
 
+// executeWithLogger runs the program with os.Stdin, piping stdout and stderr 
to bufLogger,
+// and flushes the logger at ERROR severity on failure or DEBUG severity on 
success.
+func executeWithLogger(ctx context.Context, bufLogger *tools.BufferedLogger, 
prog string, args ...string) error {
+       err := execx.ExecuteEnvWithIO(nil, os.Stdin, bufLogger, bufLogger, 
prog, args...)
+       return bufLogger.Flush(ctx, err)
+}
+
+// executeWithOutput runs the program with os.Stdin, capturing stdout in a 
byte buffer
+// while piping stderr to bufLogger, and flushes the logger at ERROR severity 
on failure or DEBUG severity on success.
+func executeWithOutput(ctx context.Context, bufLogger *tools.BufferedLogger, 
prog string, args ...string) ([]byte, error) {
+       var stdout bytes.Buffer
+       err := execx.ExecuteEnvWithIO(nil, os.Stdin, &stdout, bufLogger, prog, 
args...)
+       if flushErr := bufLogger.Flush(ctx, err); flushErr != nil {
+               return nil, flushErr
+       }
+       return stdout.Bytes(), nil
+}
+
 // pipInstallRequirements installs the given requirement, if present.
 func pipInstallRequirements(ctx context.Context, logger *tools.Logger, files 
[]string, dir, name string) error {
        pythonVersion, err := expansionx.GetPythonVersion()
@@ -62,12 +80,9 @@ func pipInstallRequirements(ctx context.Context, logger 
*tools.Logger, files []s
                        // also installs dependencies. The key is that if all 
the packages have
                        // been installed in the first round then this command 
will be a no-op.
                        args = []string{"-m", "pip", "install", "-r", 
filepath.Join(dir, name), "--no-cache-dir", "--disable-pip-version-check", 
"--find-links", dir}
-                       err := execx.ExecuteEnvWithIO(nil, os.Stdin, bufLogger, 
bufLogger, pythonVersion, args...)
-                       if err != nil {
-                               bufLogger.FlushAtError(ctx)
+                       if err := executeWithLogger(ctx, bufLogger, 
pythonVersion, args...); err != nil {
                                return fmt.Errorf("PIP failed to install 
dependencies, got %s. This error may be unrecoverable, see %s for more 
information", err, unrecoverableURL)
                        }
-                       bufLogger.FlushAtDebug(ctx)
                        return nil
                }
        }
@@ -121,23 +136,16 @@ func pipInstallPackage(ctx context.Context, logger 
*tools.Logger, files []string
                                if pipNoBuildIsolation {
                                        args = append(args, 
"--no-build-isolation")
                                }
-                               err := execx.ExecuteEnvWithIO(nil, os.Stdin, 
bufLogger, bufLogger, pythonVersion, args...)
-                               if err != nil {
-                                       bufLogger.FlushAtError(ctx)
+                               if err := executeWithLogger(ctx, bufLogger, 
pythonVersion, args...); err != nil {
                                        return fmt.Errorf("PIP failed to 
install dependencies, got %s. This error may be unrecoverable, see %s for more 
information", err, unrecoverableURL)
-                               } else {
-                                       bufLogger.FlushAtDebug(ctx)
                                }
                                args = []string{"-m", "pip", "install", 
"--no-cache-dir", "--disable-pip-version-check", filepath.Join(dir, 
packageSpec)}
                                if pipNoBuildIsolation {
                                        args = append(args, 
"--no-build-isolation")
                                }
-                               err = execx.ExecuteEnvWithIO(nil, os.Stdin, 
bufLogger, bufLogger, pythonVersion, args...)
-                               if err != nil {
-                                       bufLogger.FlushAtError(ctx)
+                               if err := executeWithLogger(ctx, bufLogger, 
pythonVersion, args...); err != nil {
                                        return fmt.Errorf("PIP failed to 
install dependencies, got %s. This error may be unrecoverable, see %s for more 
information", err, unrecoverableURL)
                                }
-                               bufLogger.FlushAtDebug(ctx)
                                return nil
                        }
 
@@ -146,12 +154,9 @@ func pipInstallPackage(ctx context.Context, logger 
*tools.Logger, files []string
                        if pipNoBuildIsolation {
                                args = append(args, "--no-build-isolation")
                        }
-                       err := execx.ExecuteEnvWithIO(nil, os.Stdin, bufLogger, 
bufLogger, pythonVersion, args...)
-                       if err != nil {
-                               bufLogger.FlushAtError(ctx)
+                       if err := executeWithLogger(ctx, bufLogger, 
pythonVersion, args...); err != nil {
                                return fmt.Errorf("PIP failed to install 
dependencies, got %s. This error may be unrecoverable, see %s for more 
information", err, unrecoverableURL)
                        }
-                       bufLogger.FlushAtDebug(ctx)
                        return nil
                }
        }

Reply via email to