Thread XML-Dateien zusammenfügen

Adriano10

Bekanntes Mitglied
Auf dem SFTP-Server liegen über 500000 Dateien, die verarbeitet werden müssen, also pro Datei müssen 25000 zusammengefügt werden.

Ich habe mich für Thread entschieden, aber es ist sehr langsam, hätte jemand Optimierungsvorschläge? Im Umgang mit Thread bin ich so zu sagen Anfänger.

Mit besten Grüßen


Java:
   Vector<ChannelSftp.LsEntry> entries = sftp.ls(SOURCE);
        Lock sftpLock = new ReentrantLock();
        int numThreads = 4;
        int filesPerThread = entries.size() / numThreads;
        int remainingFiles = entries.size() % numThreads;

        Thread thread1 = new Thread(() -> {
            int startIndex = 0;
            int endIndex = filesPerThread;
            processXML(startIndex, endIndex, entries, sftp, sftpLock, xmlReadWriter111);
        });

        Thread thread2 = new Thread(() -> {
            int startIndex = filesPerThread;
            int endIndex = startIndex + filesPerThread;
            processXML(startIndex, endIndex, entries, sftp, sftpLock, xmlReadWriter111);
        });

        Thread thread3 = new Thread(() -> {
            int startIndex = 2 * filesPerThread;
            int endIndex = startIndex + filesPerThread;
            processXML(startIndex, endIndex, entries, sftp, sftpLock, xmlReadWriter111);
        });

        Thread thread4 = new Thread(() -> {
            int startIndex = 3 * filesPerThread;
            int endIndex = startIndex + filesPerThread + remainingFiles;
            processXML(startIndex, endIndex, entries, sftp, sftpLock, xmlReadWriter111);
        });

        thread1.start();
        thread2.start();
        thread3.start();
        thread4.start();

        thread1.join();
        thread2.join();
        thread3.join();
        thread4.join();

        sftp.disconnect();
        session.disconnect();
    }


    private static void processXML(int startIndex, int endIndex, Vector<ChannelSftp.LsEntry> entries, ChannelSftp sftp,
                                   Lock sftpLock, XMLReadWriter111 xmlReadWriter111) {
        List<Document> materials = new ArrayList<>();
        List<Document> erp_mark = new ArrayList<>();
        List<Document> crossreference = new ArrayList<>();
        List<Document> dokuinfosatz = new ArrayList<>();
        for (int k = startIndex; k < endIndex; k++) {
            ChannelSftp.LsEntry entry = entries.get(k);
            //Material"
            if (!entry.getAttrs().isDir() && entry.getFilename().startsWith("Material")) {
                String remoteFile = SOURCE + entry.getFilename();
                System.out.println(Thread.currentThread().getName() + " :" + k);
                try {
                    sftpLock.lock();
                    materials.add(parseXml(sftp.get(remoteFile)));
                } catch (SftpException e) {
                    e.printStackTrace();
                } finally {
                    sftpLock.unlock();
                }
            }

            if (materials.size() == 25000 && !materials.isEmpty()) {
                try {
                    sftpLock.lock();
                    xmlReadWriter111.createOneXMLFile(materials,
                            "C:\\Users\\test\\IdeaProjects\\XMLDateien\\xml", Tags.PRODUCTS.value());
                    materials.clear();
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    sftpLock.unlock();
                }
            }
            //ERP_MARKE
            if (!entry.getAttrs().isDir() && entry.getFilename().startsWith("ERP_MARKE")) {
                String remoteFile = SOURCE + entry.getFilename();
                System.out.println("B " + k);
                try {
                    sftpLock.lock();
                    erp_mark.add(parseXml(sftp.get(remoteFile)));
                } catch (SftpException e) {
                    e.printStackTrace();
                } finally {
                    sftpLock.unlock();
                }
            }

            if (erp_mark.size() == 25000 && !erp_mark.isEmpty()) {
                try {
                    sftpLock.lock();
                    xmlReadWriter111.createOneXMLFile(erp_mark,
                            "C:\\Users\\test\\IdeaProjects\\XMLDateien\\xml", Tags.PRODUCTS.value());
                    erp_mark.clear();
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    sftpLock.unlock();
                }
            }

            //Crossreference
            if (!entry.getAttrs().isDir() && entry.getFilename().startsWith("Crossreference")) {
                String remoteFile = SOURCE + entry.getFilename();
                System.out.println("Crossreference " + k);
                try {
                    sftpLock.lock();
                    crossreference.add(parseXml(sftp.get(remoteFile)));
                } catch (SftpException e) {
                    e.printStackTrace();
                } finally {
                    sftpLock.unlock();
                }
            }

            if (crossreference.size() == 25000 && !crossreference.isEmpty()) {
                try {
                    sftpLock.lock();
                    xmlReadWriter111.createOneXMLFile(crossreference,
                            "C:\\Users\\test\\IdeaProjects\\XMLDateien\\xml", Tags.PRODUCTS.value());
                    erp_mark.clear();
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    sftpLock.unlock();
                }
            }

            //Dokuinfosatz
            if (!entry.getAttrs().isDir() && entry.getFilename().startsWith("Dokuinfosatz")) {
                String remoteFile = SOURCE + entry.getFilename();
                System.out.println("Dokuinfosatz" + k);
                try {
                    sftpLock.lock();
                    dokuinfosatz.add(parseXml(sftp.get(remoteFile)));
                } catch (SftpException e) {
                    e.printStackTrace();
                } finally {
                    sftpLock.unlock();
                }
            }

            if (dokuinfosatz.size() == 25000 && !dokuinfosatz.isEmpty()) {
                try {
                    sftpLock.lock();
                    xmlReadWriter111.createOneXMLFile(dokuinfosatz,
                            "C:\\Users\\test\\IdeaProjects\\XMLDateien\\xml", Tags.ASSETS.value());
                    erp_mark.clear();
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    sftpLock.unlock();
                }
            }
        }

        if (!materials.isEmpty()) {
            try {
                sftpLock.lock();
                xmlReadWriter111.createOneXMLFile(materials,
                        "C:\\Users\\test\\IdeaProjects\\XMLDateien\\xml", Tags.PRODUCTS.value());
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                sftpLock.unlock();
            }
        }
        if (!erp_mark.isEmpty()) {
            try {
                sftpLock.lock();
                xmlReadWriter111.createOneXMLFile(erp_mark,
                        "C:\\Users\\test\\IdeaProjects\\XMLDateien\\xml", Tags.PRODUCTS.value());
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                sftpLock.unlock();
            }
        }

        if (!crossreference.isEmpty()) {
            try {
                sftpLock.lock();
                xmlReadWriter111.createOneXMLFile(crossreference,
                        "C:\\Users\\test\\IdeaProjects\\XMLDateien\\xml", Tags.PRODUCTS.value());
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                sftpLock.unlock();
            }
        }
        if (!dokuinfosatz.isEmpty()) {
            try {
                sftpLock.lock();
                xmlReadWriter111.createOneXMLFile(dokuinfosatz,
                        "C:\\Users\\test\\IdeaProjects\\XMLDateien\\xml", Tags.ASSETS.value());
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                sftpLock.unlock();
            }
        }
    }
 
Naja, du machst ja auch viel innerhalb deiner Sperre, also solange eine Datei verarbeitet wird kann keine andere verarbeitet werden, damit bist du effektiv wieder einzelnen verarbeiten von einer Datei nach der anderen.

Du willst Dateien holen, und diese dann der Logik zufuehren. Oder, du holst die Datei und entsperrst dann dein SFTP wieder, und verarbeitest sie erst dann nachdem du das SFTP wieder freigegeben hast damit die naechste Datei schonmal geholt werden kann. Alternativ koenntest du auch mehrere SFTP Sitzungen verwenden, naemlich eine pro Thread, dann ist wirklich alles parallel.
 
Also, du machst nicht alles innerhalb der Sperre, aber recht viel. Du solltest da zumindest holen und verarbeiten trennen so dass nur das Holen innerhalb der Sperre passiert, und alles andere nicht. Aber die beste Loesung waere mehrere SFTP-Verbindungen, eine pro Thread. Weil dann kannst du die wirklich parallel verarbeiten.
 
Also, du machst nicht alles innerhalb der Sperre, aber recht viel. Du solltest da zumindest holen und verarbeiten trennen so dass nur das Holen innerhalb der Sperre passiert, und alles andere nicht. Aber die beste Loesung waere mehrere SFTP-Verbindungen, eine pro Thread. Weil dann kannst du die wirklich parallel verarbeiten.
vielen Dank, probiere ich mall dann so, also du meinst ungefähr so?
Java:
        Thread thread1 = new Thread(() -> {
            int startIndex = 0;
            int endIndex = filesPerThread;
            Session session = jsch.getSession(username, remoteHost, 22);
            processXML(startIndex, endIndex, entries, sftp, sftpLock, xmlReadWriter111);
        });

        Thread thread2 = new Thread(() -> {
            int startIndex = filesPerThread;
            int endIndex = startIndex + filesPerThread;
            Session session = jsch.getSession(username, remoteHost, 22);
            processXML(startIndex, endIndex, entries, sftp, sftpLock, xmlReadWriter111);
        });

        Thread thread3 = new Thread(() -> {
            int startIndex = 2 * filesPerThread;
            int endIndex = startIndex + filesPerThread;
            Session session = jsch.getSession(username, remoteHost, 22);
            processXML(startIndex, endIndex, entries, sftp, sftpLock, xmlReadWriter111);
        });

        Thread thread4 = new Thread(() -> {
            int startIndex = 3 * filesPerThread;
            int endIndex = startIndex + filesPerThread + remainingFiles;
            Session session = jsch.getSession(username, remoteHost, 22);
            processXML(startIndex, endIndex, entries, sftp, sftpLock, xmlReadWriter111);
        });
 
Also, du machst nicht alles innerhalb der Sperre, aber recht viel. Du solltest da zumindest holen und verarbeiten trennen so dass nur das Holen innerhalb der Sperre passiert, und alles andere nicht. Aber die beste Loesung waere mehrere SFTP-Verbindungen, eine pro Thread. Weil dann kannst du die wirklich parallel verarbeiten.
also sftpLock.lock(); and sftpLock.unlock(); ist nicht mehr nötig oder?
 
das ist gut aber, threads müssen auf gleiche Ressourcen in der Methode processXML zugreifen und daher ohne synchronisiere gibt es Lesekonflikt
Welche denn? Du definierst da ein paar Listen, aber keine davon werden auf irgendeine Art und Weise geteilt (sind nur in der Methode)?

Wenn dem so ist, musst du dir dann entweder eine feinere Zugriffskontrolle ueberlegen (so wenig Sperren wie moeglich), oder, ich weisz jetzt nicht wie genau deine Anforderung aussieht, du findest einen Weg jeden Thread fuer sich arbeiten zu lassen und dann die Ergebnisse zussamen zu fuehren.
 
Mir fällt es schwer, den ganzen Ansatz so nachzuvollziehen.

Ich bin ein Freund davon, dass man Code einfach hält. Du hast ein Vector mit den Einträgen - warum nicht einfach alles per parallelStream verarbeiten? Dann sparst Du Dir viel Code rund um Anzahl der Threads und deren Erzeugung. Die Java Runtime achtet dann darauf, dass du eine relativ optimale Anzahl Threads bekommst. Also bei vielen Cores hast Du mehr aktive Threads und so. Und ganz wichtig: Der Code wird einfacher!

Und dann wirklich die große Bitte: Schreibe sauberen Code!

a) doppelter Code! Du hast eine Variable mit Anzahl Threads und jeder Thread wird gleich erzeugt. Das schreit doch förmlich nach einer Schleife statt einfach nur Code mit Copy & Paste zu kopieren.

b) () -> { Wieso packst Du das aus dem Block nicht einfach in eine Methode und rufst diese dann auf? Lambda Expressions mit einem Block sind immer der Grund, wieso ich meine, dass statische Codeanalysetools eine tens unit brauchen um Entwicklern Stromstöße zu geben 🙂 (Sorry, ist natürlich nur ein Spaß. Aber so manche Dinge machen Code wirklich nur unleserlich. Den Spaß bitte daher etwas verzeihen)

c) processXML - der Code geht von Zeile 46 bis 207. Muss ich da mehr zu schreiben?

Das wirkt sehr undurchdacht. Und ich habe da gewisse Probleme, die Motivation zu finden um da weiter drüber nach zu denken.


Dann evtl. noch eine Art minimales Connection Pooling. Die SSH Session würde ich also cachen. Dann kann jeder Aufruf einfach auf dem cache ein getSession() aufrufen und bekommt dann eine Instanz, die Autoclosable implementiert und eine Session kapselt. nach verwendung wird auf der Instanz ein close() aufgerufen so dass die Session für den nächsten Thread bereit steht. Ist relativ einfach und schnell umsetzbar denke ich.


Aber natürlich auch ganz wichtig: Was @Robert Zenz angesprochen hatte: Die Datenstruktur muss stimmig sein. Du hast irgend ein Aufbau mit
Java:
        List<Document> materials = new ArrayList<>();
        List<Document> erp_mark = new ArrayList<>();
        List<Document> crossreference = new ArrayList<>();
        List<Document> dokuinfosatz = new ArrayList<>();
Das scheinen irgendwelche zentralen Daten zu sein, die Du aufbauen willst. Das gehört in eine Klasse. Natürlich incl. der notwendigen Logik, die man da so braucht um Dinge hinzu zu fügen. Da wir über das fachliche nicht Bescheid wissen, können wir natürlich nichts sagen. Aber evtl. gibt es von einer Datei eine Referenz auf eine andere Datei, die noch nicht verarbeitet wurde. So Dinge muss man natürlich abdecken. Sowas will man dann ggf. einfügen incl. einem Vermerk, dass da noch die Basis fehlt. Dann hat man am Ende ggf. eine Liste mit Referenzen, die ins "Leere" zeigen oder so ... Aber das ist halt ein anderes Thema. Das ist etwas, das dann - sauber unterteilt in Klassen / Methoden erstellt werden muss.

Sowas im idealen Fall übrigens incl. UnitTests, die so fachliche Anforderungen auch testen. (Und die Unit Tests sind dann ein super Ort um die fachlichen Vorgaben zu dokumentieren mit JavaDoc Kommentaren 🙂 )

Das einfach nur als ein paar allgemeine Informationen und Hinweise. Auch wenn dies womöglich nicht das direkte Kernproblem ansprechen, das Du lösen willst, hoffe ich doch sehr, dass diese allgemeinen Hinweisen dazu führen, dass Du besseren Code schreiben kannst, den Du dann auch einfacher selbst verstehen und ggf. anpassen kannst.
 
Mir fällt es schwer, den ganzen Ansatz so nachzuvollziehen.

Ich bin ein Freund davon, dass man Code einfach hält. Du hast ein Vector mit den Einträgen - warum nicht einfach alles per parallelStream verarbeiten? Dann sparst Du Dir viel Code rund um Anzahl der Threads und deren Erzeugung. Die Java Runtime achtet dann darauf, dass du eine relativ optimale Anzahl Threads bekommst. Also bei vielen Cores hast Du mehr aktive Threads und so. Und ganz wichtig: Der Code wird einfacher!

Und dann wirklich die große Bitte: Schreibe sauberen Code!

a) doppelter Code! Du hast eine Variable mit Anzahl Threads und jeder Thread wird gleich erzeugt. Das schreit doch förmlich nach einer Schleife statt einfach nur Code mit Copy & Paste zu kopieren.

b) () -> { Wieso packst Du das aus dem Block nicht einfach in eine Methode und rufst diese dann auf? Lambda Expressions mit einem Block sind immer der Grund, wieso ich meine, dass statische Codeanalysetools eine tens unit brauchen um Entwicklern Stromstöße zu geben 🙂 (Sorry, ist natürlich nur ein Spaß. Aber so manche Dinge machen Code wirklich nur unleserlich. Den Spaß bitte daher etwas verzeihen)

c) processXML - der Code geht von Zeile 46 bis 207. Muss ich da mehr zu schreiben?

Das wirkt sehr undurchdacht. Und ich habe da gewisse Probleme, die Motivation zu finden um da weiter drüber nach zu denken.


Dann evtl. noch eine Art minimales Connection Pooling. Die SSH Session würde ich also cachen. Dann kann jeder Aufruf einfach auf dem cache ein getSession() aufrufen und bekommt dann eine Instanz, die Autoclosable implementiert und eine Session kapselt. nach verwendung wird auf der Instanz ein close() aufgerufen so dass die Session für den nächsten Thread bereit steht. Ist relativ einfach und schnell umsetzbar denke ich.


Aber natürlich auch ganz wichtig: Was @Robert Zenz angesprochen hatte: Die Datenstruktur muss stimmig sein. Du hast irgend ein Aufbau mit
Java:
        List<Document> materials = new ArrayList<>();
        List<Document> erp_mark = new ArrayList<>();
        List<Document> crossreference = new ArrayList<>();
        List<Document> dokuinfosatz = new ArrayList<>();
Das scheinen irgendwelche zentralen Daten zu sein, die Du aufbauen willst. Das gehört in eine Klasse. Natürlich incl. der notwendigen Logik, die man da so braucht um Dinge hinzu zu fügen. Da wir über das fachliche nicht Bescheid wissen, können wir natürlich nichts sagen. Aber evtl. gibt es von einer Datei eine Referenz auf eine andere Datei, die noch nicht verarbeitet wurde. So Dinge muss man natürlich abdecken. Sowas will man dann ggf. einfügen incl. einem Vermerk, dass da noch die Basis fehlt. Dann hat man am Ende ggf. eine Liste mit Referenzen, die ins "Leere" zeigen oder so ... Aber das ist halt ein anderes Thema. Das ist etwas, das dann - sauber unterteilt in Klassen / Methoden erstellt werden muss.

Sowas im idealen Fall übrigens incl. UnitTests, die so fachliche Anforderungen auch testen. (Und die Unit Tests sind dann ein super Ort um die fachlichen Vorgaben zu dokumentieren mit JavaDoc Kommentaren 🙂 )

Das einfach nur als ein paar allgemeine Informationen und Hinweise. Auch wenn dies womöglich nicht das direkte Kernproblem ansprechen, das Du lösen willst, hoffe ich doch sehr, dass diese allgemeinen Hinweisen dazu führen, dass Du besseren Code schreiben kannst, den Du dann auch einfacher selbst verstehen und ggf. anpassen kannst.
hier gebe ich vollkommen Recht, ich achte immer auch möglicherweise auf Codequalität, aber das ist keine fertiges Code, ich hatte nur angefangen, wenn es läuft, danach mache ich immer Verbesserung
 
hier gebe ich vollkommen Recht, ich achte immer auch möglicherweise auf Codequalität, aber das ist keine fertiges Code, ich hatte nur angefangen, wenn es läuft, danach mache ich immer Verbesserung
Das ist etwas, das durchaus normal ist. Aber meine Erfahrung hat recht deutlich gezeigt, dass es deutlich besser ist, wenn man von Anfang an entsprechend strukturiert. Alleine schon, um auch von Anfang an mit Unit Tests zu arbeiten.

Klar, wenn man etwas ausprobiert, dann sind die Unit Tests keine wirklichen Unit Tests. Dann greift man auf einen SSH Server zu und arbeitet da mit Dateien oder so. Aber dennoch strukturiert man alles etwas.

Gerade bei so Dingen, die schnell etwas komplexer werden, ist man dann unter dem Strich deutlich schneller unterwegs.

Und man hat immer das Risiko, dass Du so eine Methode hast, die erst einmal macht, was sie soll. Das Umschreiben dazu fehlt die Zeit und schon hast Du einen Code wie das processXML in Produktion. Das ist dann eine (annähernd) ungetestete Methode und es ist nur eine Frage der Zeit, bis es notwendig wird, da Änderungen dran vorzunehmen. Oder in 1 Jahr oder so mal überlegen, was die genauen Spezifikationen waren, die da umgesetzt wurden 🙂

(Und wie wichtig so Dinge sind, merkst Du, wenn Du in einem Projekt bist, das seit > 30 Jahren läuft. Das geht nur mit sehr hoher Code Qualität. Einiges ist aus heutiger Sicht nicht gut, aber die damalige, hohe Qualität hat dazu geführt, dass die Software auch nach all diesen Jahrzehnten noch immer angepasst werden kann und immer noch stabil im Betrieb ist. Und das ohne wirkliche Ausfälle in Produktion - was wichtig ist, denn an gewisser Software hängt ggf. auch mit die Existenz eines Konzerns 🙂 )
 

Neue Themen


Zurück
Oben