使用Java+Spark Streaming读取日志文件无输出的问题求助
Hey there, let's troubleshoot why your Java + Spark streaming setup is only spitting out repeated content instead of your logback-generated logs. I’ve run into this exact issue before, so here’s a breakdown of what’s going wrong and how to fix it:
That repeated output is almost certainly Spark’s default streaming monitoring logs—your actual log file content isn’t being read at all. The most common culprits are:
- You’re using the old
textFileStreamAPI, which doesn’t handle continuously updated single files (Logback writes to one file until it rolls over, and Spark won’t pick up incremental changes here). - Your file path configuration is off, or Spark doesn’t have read access to
c:/test. - You’re not using the right output mode, so Spark is reprinting the same empty/metadata batches over and over.
The solution is to use Spark Structured Streaming (the modern, more reliable API for streaming workloads) which handles continuously updated files much better. Here’s a working Java example tailored to your Logback log file:
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.streaming.StreamingQuery; import org.apache.spark.sql.streaming.StreamingQueryException; public class LogStreamProcessor { public static void main(String[] args) throws StreamingQueryException { // Initialize SparkSession (local mode for testing; remove master() in production) SparkSession spark = SparkSession.builder() .appName("LogbackStreamProcessor") .master("local[*]") .getOrCreate(); // Read from the log directory (Spark monitors all files here for changes) Dataset<Row> logStream = spark.readStream() .format("text") .option("maxFilesPerTrigger", 1) // Process only new content per batch .option("cleanSource", "archive") // Optional: Archive processed files to avoid re-reading .option("sourceArchiveDir", "c:/test/processed_logs") // Create this directory first! .load("c:/test"); // Print logs to console (your temporary requirement) StreamingQuery consoleQuery = logStream.writeStream() .outputMode("append") // Only print new, unprocessed log lines .format("console") .option("truncate", false) // Don't cut off long log lines .start(); consoleQuery.awaitTermination(); } }
Key Notes for This Code:
- Point to the directory, not the single log file: Spark monitors the entire
c:/testfolder—if you configure Logback to roll logs into new files (by size/time), Spark will automatically pick up each new rollover file. appendoutput mode: Ensures you only see new log lines each batch, not repeated content.- Archive option: Prevents Spark from reprocessing the same log files over and over by moving them to
c:/test/processed_logsafter handling.
Once your console output is working, switching to writing logs to a database is straightforward. Here’s an example for MySQL (adjust for your database of choice):
// Replace the console writeStream block with this database code StreamingQuery dbQuery = logStream.writeStream() .outputMode("append") .format("jdbc") .option("url", "jdbc:mysql://your-db-host:3306/your_database") .option("dbtable", "application_logs") // Your target table name .option("user", "db_username") .option("password", "db_password") .option("checkpointLocation", "c:/test/spark_checkpoint") // Critical for fault tolerance .start(); dbQuery.awaitTermination();
Pre-requisite: Set Up Your Database Table
Create a table to store the logs first—here’s a simple MySQL schema:
CREATE TABLE application_logs ( id INT AUTO_INCREMENT PRIMARY KEY, log_content TEXT NOT NULL, ingested_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP );
- Verify Spark has read/write permissions for
c:/testand the archive/checkpoint directories. - Check Spark’s own logs (usually in
spark/logsif running locally) for any file-access errors. - If you’re stuck on the old DStream API: Don’t use
textFileStreamfor continuously updated single files—it’s designed for one-time file additions, not ongoing writes.
内容的提问来源于stack exchange,提问作者Sam

