Wednesday, May 27, 2026

Optimizing Sequential File Processing in Databricks Medallion Architecture

 

The Databricks job processes data on a file-by-file basis. Based on our analysis of daily file arrival trends, the maximum number of files processed in a single day is 590, while the minimum number of files processed in a day is 194.

Files are received in a non-uniform and unpredictable manner, with no consistent batching pattern. Due to this irregular arrival pattern, it is challenging to accurately define a fixed schedule or frequency for job execution.

File sizes vary considerably, with some files containing more than 516 rows. These larger file sizes were not considered during the earlier performance analysis.

Currently, our workloads are executed on a Databricks single-node cluster configured with Standard_D8ds_v5 (8 cores, 32 GB RAM). We have observed that this setup is resulting in longer processing times.

 After optimizing the Databricks code to enable parallel processing with up to 8 concurrent threads, we conducted performance testing.

The results showed that processing 150 files, each with an average size of approximately 1 KB, takes 18 minutes and 36 seconds in total. This equates to an average processing time of approximately 7.4 seconds per file.

16.4 LTS (includes Apache Spark 3.5.2, Scala 2.12) Standard_D8ds_v5 32GB 8 Cores and that can roughly cost as much as follows –

🔹 D8ds_v5 cluster

  • DBU: $2,376
  • VM: ~$330 👉 Total ≈ $2,700/month 

While higher-capacity cluster configurations can significantly improve file processing performance, they also incur increased costs. The following cluster configuration, leveraging 16 concurrent threads, can process 150 files considerably faster compared to the current setup. Time may vary slightly depending on file size.

 16.4 LTS (includes Apache Spark 3.5.2, Scala 2.12) Standard_D16ds_v5 64GB 16 Cores roughly cost as much as follows -

🔹 D16ds_v5 cluster

  • DBU: $4,752
  • VM: ~$650 👉 Total ≈ $5,400/month

 At code level following changes are done:

The process consists of two primary components: file validation and data validation. The file validation stage verifies key attributes such as file name, file length, extension, and sequence number. As the sequence integrity must be maintained, it becomes challenging to process files in parallel during this stage.

The data validation stage is responsible for validating the data and inserting it into the database. This phase can be executed in parallel; accordingly, we have implemented multithreading to enable concurrent processing and improve performance.

At this stage, we observed that sequential file processing was the most time-consuming compared to the other two layers. Upon further investigation and identifying additional optimization opportunities, we executed sequential file validation in parallel based on categories, which resulted in significant performance improvement.

In a multithreaded environment, the number of available CPU cores plays a critical role in performance. The number of concurrent threads can be proportionally adjusted based on the core capacity of the Databricks cluster; however, this also leads to a corresponding increase in cost.

Key Point: Optimizing queries and code is critical. Identifying further opportunities for performance improvement should be treated as an ongoing, continuous process. Avoid increasing resources unnecessarily, as this leads to higher Databricks cluster costs.

Wednesday, May 20, 2026

Databricks-Centric Lakehouse Architecture

 Case Study 

Existing Application to a Databricks-Centric Lakehouse Platform

Executive Summary

The current application architecture relies on a multi-layered Azure-based setup involving React App Service, Azure Functions, and Cosmos DB for data ingestion, processing, and visualization. While functional, this design introduces unnecessary complexity, operational overhead, and multiple points of dependency.

This proposal outlines a modernized, simplified, and scalable architecture leveraging Databricks Apps, Unity Catalog, Delta Lake, and Databricks AI (Genie) to streamline the system into a unified platform. The proposed approach reduces service dependencies, enhances governance, improves performance, and enables native AI-driven insights.

Current Architecture Overview

The existing solution consists of the following components:

  • Frontend Application (React JS) hosted on Azure App Service
  • Azure Functions acting as middleware for:
    • Communication with Databricks / ADLS
    • Communication with Cosmos DB
    • Notification handling
  • Cosmos DB used for:
    • Data ingestion and storage
    • Querying for visualization
  • Databricks / ADLS accessed indirectly via Azure Functions

Key Functional Capabilities

  • Capture user inputs from frontend
  • Update JSON payloads
  • Processed Input data based on updated JSON
  • Ingest data into Cosmos DB
  • Retrieve data for visualization
  • Send user notifications

Challenges with Current Architecture

The current design introduces several limitations:

  • High Dependency Chain
    • Tight coupling between frontend, Azure Functions, Cosmos DB, and Databricks
  • Operational Complexity
    • Multiple services to maintain and monitor
    • Increased DevOps overhead
  • Performance Overhead
    • Multiple network hops between services
    • Increased latency for data access
  • Governance Fragmentation
    • Data access control spread across services
    • Limited centralized governance
  • Limited AI Enablement
      • Minimal integration with advanced analytics and AI capabilities

Proposed Architecture

Proposed transitioning to a Databricks-centric unified architecture that consolidates application, data, and AI capabilities into a single platform.

Core Components

 Databricks Apps (Frontend Layer)

  • Replace Azure App Service + Azure Functions
  • Provide UI for:
    • Capturing user input
    • Managing JSON data
    • Rendering visualizations

Delta Lake on ADLS (Data Layer)

  • Replace/augment Cosmos DB
  • Store:
    • Structured data (Delta tables)
    • Semi-structured JSON data 

 Unity Catalog (Governance Layer)

  • Centralized control for:
    • Data access (RBAC/ABAC)
    • Data lineage
    • Security policies 

 Databricks SQL Warehouse (Query Engine)

  • High-performance query execution
  • Enables dashboards and app-driven queries

 Databricks AI / Genie (Optional Layer)

  • Natural language querying (NL → SQL)
  • AI-driven insights and summarization

 Databricks Dashboards

  • Replace custom-coded visualization logic
  • Provide governed, reusable visual reporting

Proposed Functional Flow

User → Databricks App (SSO via Entra ID)

     → Direct interaction with Delta Tables (via SQL Warehouse)

     → Unity Catalog enforces access controls

     → Data stored/retrieved from ADLS (Delta + JSON)

     → Visualization via built-in dashboards or app UI

     → Optional: AI-driven insights via Genie

Key Improvements

 1. Reduced Dependency Footprint

  • Eliminates:
    • Azure Functions
    • Intermediate API layers
  • Reduces system complexity

 2. Unified Data Platform

  • Single platform for:
    • Data ingestion
    • Storage
    • Processing
    • Visualization
    • AI

 3. Enhanced Governance

  • Centralized through Unity Catalog:
    • Fine-grained access control
    • Auditability
    • Data lineage

 4. Improved Performance

  • Direct data access (no intermediaries)
  • Optimized query execution via SQL Warehouse
  • Reduced network overhead

 5. Cost Optimization

  • Elimination of:
    • Cosmos DB RU provisioning
    • Azure Function execution costs
  • Pay-per-use model with serverless compute

 6. Native AI Enablement

  • Use Databricks Genie to:
    • Enable natural language interactions
    • Generate insights without manual queries
  • Reduce need for custom analytics logic

7. Simplified Visualization Strategy

  • Replace custom graph rendering with:
    • Databricks Dashboards (no/low code)
  • Maintain flexibility via:
    • Optional custom visualization (Plotly/Streamlit)

Expected Outcomes

Area

Impact

Architecture complexity

⬇ Reduced significantly

Performance

⬆ Improved

Cost

⬇ Optimized (20–50%)

Governance

⬆ Centralized

Maintainability

⬆ Simplified

AI capability

⬆ Enable


        Key Points:

            Centralized Secure Architecture:

                        A unified, identity-driven Databricks centric security framework enables proactive                                    risk management and business scalability.

            Security as Business Enabler:

                        Transforming security from a reactive role into a proactive driver of trust and growth                                 supports business objectives.

            Investment in Resilience:

                        A modern, scalable, and AI-enabled Lakehouse solution is an investment yielding long-                            term confidence and sustainable success.

High-level design

Trade-offs / Considerations

Area

Consideration

Frontend flexibility

Databricks Apps less mature than full React

Cosmos DB

Keep only if low-latency transactional workloads needed

Skill shift

Teams need Databricks-centric skills

Vendor lock-in

More reliance on Databricks ecosystem


This workload is data engineering (ingestion + validation → gold), NOT an app / BI / interactive querying workload.  Therefore, we cannot completely avoid compute (Databricks cluster or equivalent) But we can replace traditional clusters with more cost-efficient options.

Databricks Notebook fetch data from ADLS and Event Hub, therefore ingestion layer and process layer is separate and independent. Data can be ingested continuously, and Databricks job can run batch by batch in a day from Monday to Friday, if business permits and that can reduce cluster cost. Second option is to use job cluster; however, job cluster will take time to initialize and installed libraries to be ready to process.

Conclusion

The proposed re-architecture transforms the current system into a modern, scalable, and AI-enabled Lakehouse solution. By consolidating multiple services into Databricks, the organization can achieve:

  • Reduced operational overhead
  • Improved performance and scalability
  • Stronger data governance
  • Enhanced user experience with built-in AI capabilities

This approach aligns with enterprise best practices for data platforms and provides a future-ready foundation for advanced analytics and intelligent applications.

Tuesday, July 20, 2021

Upsert Parquet Data Incrementally

Incremental data load is very easy now a days. Next generation Databricks Delta allows us to upsert and delete records efficiently in data lakes. However, it's a bit tedious to emulate a function that can upsert parquet table incrementally like Delta. In this article I'm going to throw some light on the subject.

Hadoop follow WORM (write once and read multiple time) that doesn't allow us to delete rows from Data Frame. But then question appears how to handle restatement data? When last week data got changed in current week, we need to update the row with latest value in master table.

In this example we'll process historical and latest data before overwriting existing table, however, for large data sets, it will impact engine performance. We need to segregate input data in such a way so that only the partition gets update where there is a change. For that, data wrangling is most important.

In below set of example Data Frame, assume Table 1 is our previous week data and Table tow has been received in current week. Significant point is that the row of week ID 260 (Sales value) got changed in Table 2. We have to keep that change in master data.

Let's prepare three sets of test Data Frame and apply upsert function with first two. Using "testthat" R library we'll compare the result with third Data Frame.

Sunday, October 25, 2020

Spark Processing - Leveraging Databricks Jobs

 

Running spark application in Databricks require many architectural considerations. From the beginning of choosing right cluster up to coding is million-dollar question. Not only that, post implementation, monitoring job performance and optimizing ETL jobs is another continuous process of improvement.

Here, we’ll discuss a few points that can boost up job performance and report your business at earliest.

Databricks offer two types of clusters comprising different runtime for different workload. Choosing Databricks runtime based on working area and domain is first and foremost important point to be considered.

Next point is to select worker and driver type. Before going into the depth of different types of worker and Driver, lets have a look into the function of them.

Driver:  The Driver is one of the nodes in the Cluster. The driver does not run computations, it plays the role of a master node in the Spark cluster. When you join multiple portion of Dataset from different executor, the whole data is sent to the Driver.

Worker: Workers run the Spark executors and other services required for the proper functioning of the clusters. Process of distributed workload happens on workers. Databricks runs one executor per worker node; therefore, the terms executor and worker are used interchangeably. Executors are JVMs that run on Worker nodes. These are the JVMs that run Tasks on data Partitions.

There are two types of cluster modes. Standard and high concurrency. High concurrency provides resource utilization, isolation for each notebook by creating a new environment for each one, security and sharing by multiple concurrently active users. Sharing is accomplished by pre-empting tasks to enforce fair sharing between different users. Pre-emption is configurable.

Now it would be easy to understand the requirement and select worker and driver type. Not like that, we need to consider price offered by different service. Databricks is nothing but a PaaS. It’s depends on two major instance provider - AWS and Azure. Go to the product price page of Databricks, it will offer you to select any one of two. Below are the prices offered (old one) for Microsoft Azure. (Please check the latest offer)



Now, we are bit serious about selecting the platform to run our ETL jobs. Databricks and Azure both are well documented however it is fragmented, therefore, above information will help understanding the concept quickly.

Hope we have already selected near to perfect platform based on our work type and budget. Next to choose language based on our convenient. Databricks supports multiple languages but we’ll always get the best performance with JVM-based languages like Spark-SQL, java, Scala. On top of that Apache Spark is written in Scala, therefore writing ETL in Scala will be advantageous indeed. However, every language has its own advantage, like python is bit popular, where R will be best use for plotting. Along with performance, depending on capability and availability of resource we should select language. 

Next point comes in my mind is different file formats for data stored in Apache Hadoop—including CSV, JSON, Apache Avro, and Apache Parquet. Text processing (CSV and JSON) are replaced by most people with Avro and Parquet as the main contenders. General observation of Databricks jobs reveals that when we process PARQUET format of file(read/write), at the time of shuffling number of partitions get increased compared to CSV, distribution of task and parallelism also seems to be more optimized comparatively.


When it comes to choosing Hadoop file format, there are many factors involved—such as integrating with third-party applications, schema evolution requirements, data type availability, and performance. But if performance matters, benchmarking show that Parquet would be the format to choose. However, Databricks Delta extends Apache Spark to simplify data reliability and boost Spark's performance.

Apart from that autoscaling and Databricks pools can improve performance of spark jobs, however cost involve with that. Databricks does not charge DBUs while instances are idle in the pool. Instance provider billing does apply.

Now code optimization may the last option to boost up performance. I’ll discuss the same in a separate thread here only.


Saturday, November 23, 2019

Using gRPC Client in CI/CD Pipeline

To expedite delivery almost every project has merged their Development and Operational activities under a common pipeline. The philosophy has been implemented in different way by keeping DevOps principles intact across the globe. Simplicity and Security is one of the most important aspect of DevOps Principals. It reduces operational overhead and improves Return of Investment.

Keeping that in mind I preferred to use GitLab which provides everything that is essentials for DevOps lifecycle. Every commit gets tested rigorously, then Build image and deploy in Docker Swarm or in Kubernetes cluster. GitLab delivers every feature very quickly to the end-user.There are various ways to achieve that. Continuous deployment usually gets configured through webhooks. A common pattern is to run http server locally that listens incoming HTTP requests from repository and triggers specific deployment command on every push.

Instead of using BaseHTTPRequestHandler to listen incoming request I used gRPC to call function defined in server side. gRPC claims 7 times faster that REST when receiving data & roughly 10 times faster than REST when sending data for a specific payload.It has many other advantages, like

·         Easy to understand.
·         Loose coupling between clients/server makes changes easy.
·         Low latency, highly scale-able, distributed systems.
·         language independent.

Enough talking, lets jump into the actual implementation.
The design is very simple,
  •         Run a gRPC server inside the controller system. 
  •        Write a simple bash/shell script comprising specific deployment command  
  •        Run gRPC client in pipeline on every push

Let’s define a hook, the actual method that will be called remotely by gRPC client.

My gRPC hook:



gRPC uses Protocol Buffers as the interface description language, and provides features such as authentication, bidirectional streaming and flow control,blocking or non blocking bindings, and cancellation and timeouts.Protocol Buffers is the default serialization format for sending data between clients and servers.
Lets define a protobuff :

Using above description grps_tools will generate two classes _pb2_grpc.py and _pb2.py. Run the following command and generate gRPC classes.



$ python -m grpc_tools.protoc -I. --python_out=. --grpc_python_out=. glb_hook.proto

Along with simplicity we must consider security. Therefore, generate self-signed certificate.


$ openssl req -newkey rsa:2048 -nodes -keyoutserver.key -x509 -days 365 -out server.crt


For more details visit gRPC documentation page and generate Server & client code.

Create deployment script and put that in same directory where gRPC server run


Configure client inside CI/CD to call deployment script remotely on every push. Once pipeline is ready, start server 

$ nohup python server.py &

At this point, every commit and push into the repository, pipeline will execute jobs.


Once client calls the function (that is defined inside the server) with a valid string value (command), a new process will be opened at server side, which is in this case, the script test.sh. If we have a look in to the server, we found Docker command pulling latest image and executing them in detach mode on a specific port. Once deployment is done, hit the URL and we'll see the change.



Tuesday, November 12, 2019

IoT


Studying Medical Data and analyzing them is Healthcare analytics which is improving human life span potentially by predicting various sign of diseases in advance. To protect our life, we do maintain quality of every intake. At the same time, we gone through different medical test periodically to check performance of inner system. Analysis excretion is among one of them. Excretion is the process that removes the waste products of digestion and metabolism from the body. It gets rid of by-products that the body is unable to use, many of which are toxic and incompatible with life. Data analysis of human body excretion gives various preventive medical information of an individual. A simple urinalysis is one way to find certain illnesses like Kidney diseases, Liver problem, Diabetes etc. in their earlier stages.

This article is not about the possibilities of capturing metrics by testing samples, rather than how we can make this test done automatically and alert individuals. Smart Sanitary System (sCube) is one way to accomplish that. Leave or release your body excretion publicly or privately, smart device will analysis that and report you.

Healthcare is one of the most important criteria for Smart City. Without proper health treatment and medication, a city will never be able to survive as a smart city. with the help of IoT it is possible to help people live smartly. Energy release by human body, weight gain or loss periodically, walking step analysis etc. can be done with the help of AI and IoT. There are potential opportunities for Health & Life insurance companies to serve their customer better and run the business more accurately.

For example, analysis done by the Iris iQ200 ELITE (İris Diagnostics, USA), Dirui FUS-200 (DIRUI Industrial Co., China) says that the degree of concordance between the two instruments was better than the degree of concordance between the manual microscopic method and the individual devices. Therefore, if sample can be analyzed automatically by instrument, collecting data and monitoring them would not be a challenge.

Site Reliability Engineering


It is assumed that DevOps philosophy has been adopted by every project at their own way. True implementation of DevOps is hidden in SRE - Site Reliability Engineering.

It seems every organization has its own SRE team in a fragmented form. Whenever there is an issue, we all jump into that and bring the business on track as per SLA. SRE talks about another two layers - SLI and SLO, which can be used as a filter of SLA. At any point of time, a particular matrix says Yes or No about system Health. These are all Service Level Indicators. Bindings targets of SLI is SLO. It never promises 100% availability of the site. Based on all these SLOs, Service Level Agreements are prepared transparently.

Transparently, because it accepts expectable risk – amount of failure we can have within our SLO. It is near to impossible to assure 100% availability, even if we provide service through our own fiber network, backbone and customized secure software. Due to least reliable component in the system we can grantee 100% availability all the time. Error Budget clearly shows minimum permissible loss beforehand. SRE expects failure is normal and determine how much failure we can tolerate. Error Budget helps to decide whether delivering new product quickly is important or Releasing reliable product/feature is our prime goal.

It has perfectly defined perhaps intended how to avoid Toil or operational overhead by discarding manual task so far possible. Manual, repetitive, automatable, tactical and devoid of long-term value are the characteristics of Overhead. Working manually by sitting in front of computer is not an intelligent decision. At the same time, investing 20Hrs to automate a single task which supposed to be done manually once in a month within 20 min, is not a wise idea, either.

Altogether, it seems, latter the service organization adopt SRE, sooner it will disappear from the market. Therefore, every organization should have a defined framework/model of SRE, if nothing as such is ready!! Experts says SRE is the class that implements the interface of DevOps. Case study on existing DevOps projects and implementing SRE on that can be represented as a POC.

SCM and evolution of DevOps


Software Configuration Management (SCM) is the application of Configuration Management (CM) principles in the context of Software Engineering (SE) in Computer Science (CS).

Software Configuration Management (SCM)  identifies and tracks the configuration of artifacts at various points in time and performs systematic control of changes to the configuration of artifacts for the purpose of maintaining integrity and traceability throughout the whole software life cycle.
It is essential for every project and applicable for every methodology of project Governance. To balance the demand of rapid software delivery and leveraging ROI, organizations started implementing automation everywhere. Domain and disciplines defined in SCM are categorized. Tools for every domain or for multiple disciplines were developed. All are interconnected as a Tool-chain that is expected to expedite software delivery.

In the era of rapid development and delivery, first we introduced Continuous Integration, though, before that, parallel development, source code merging, standardizing code commit by different hooks were in place. Many projects used customized script to integrate different software modules and packaged them together. A typical three team, three tire enterprise structure. Along with rapidness, automation also helps continuous improvement which is obvious. Over maturity of continuous integration (CI) we thought about Continuous Delivery (CD). Every project started talking about CI/CD but a few of them really succeed to Deliver Continuously. Integration of different software module is purely technical, however delivery involves many nontechnical, functional and client centric project related activities that demands agility, therefore, seems difficult to be continuous, so far monolithic application is concerned.

Merging Development, Operation and System Administration with the help of Automation introduced new concept called DevOps. The name reveals itself, Development+Operation, other than that there is no state forwards definition, rules, policies that guide us implementing the concept uniformly. Off course it has a principal, perhaps the goal to Developed as per requirement, Integrate, Deploy, Test automatically and getting Feedback for farther Improvement, then Release and finally Monitor for continuous operation. In background every change gets registered and ensure possibilities of rollback at any time. It seems like ITIL specified best practices of Service Lifecycle, however it is not a sequential framework, rather than it is Agile – an interactive approach where MVPs are passed through the life-cycle and delivered in a very short period. Altogether it’s an operation pipeline that accomplish the goal throughout a set of toolchains. 

Managing Automation and agility at the same time is not about simply letting loose a stream. To overcome the hurdle, monolithic application splits up into microservices. On the other hand, infrastructure becomes concise as a form of container. In combination of Microservices and Containerization organizations experiencing proven benefit of DevOps. Orchestration of huge containers is not a big hurdle today. Security and segregation mechanism imposed inside orchestration framework are simplifying coding complexity.

Principals of SCM are still being maintained silently inside DevOps engineering. Version Identification, Version Control, Artifact Versioning and Issue Tracking are essential disciplines of every software projects. Evolution of DevOps is continuing. Elimination or integration of tools name it differently, however core concept is automating SCM and IM to deliver rigorously Tested product in agile way. 

Sunday, July 20, 2014

Monitor SVN commit using Perl script.

Usually it should be accomplished by configuring svn hooks, however in some cases when we do not have svn server access we are unable to notice if developers are committing as per SCM guideline.  Rather monitoring manually it’s better to automate the process which can check all commits in a day and notify the user/administrator if the commit in improper.

That is the purpose of the script. We use the same user ID and corresponding credential that is being used by Jenkins. To notify user or administrator the script will use SMTP. At the end the script will be scheduled as a user’s cron job at 22:30 PM every day.

Let’s have a look into the script…………

#!/usr/bin/perl

use strict;
use warnings;
use Net::SMTP;
use XML::Simple;
my $release = shift(@ARGV);
my $dt = `date +%Y-%m-%d`;
my $base = 'https://your.host.com/base/url';
chomp($dt);
my %hash;
my %nuhash;
my %users;
my $details;
my $xml = new XML::Simple;
my $contain = $xml->XMLin("/path/to/jenkins-home/hudson.scm.SubversionSCM.xml");
my $user = $contain->{credentials}->{'entry'}->[0]->{'hudson.scm.SubversionSCM_\
                                                      -DescriptorImpl_-PasswordCredential'}->{'userName'} . "\n";
my $auth = $contain->{credentials}->{'entry'}->[0]->{'hudson.scm.SubversionSCM_\
                                                      -DescriptorImpl_-PasswordCredential'}->{'password'};
my $encode = `echo $auth | python -m base64 -d`;
chomp($user);
chomp($encode);
if ($release =~ /REL_VER_/){
foreach my $repo ('app1', 'app2', 'app3') {
my $repo = "$repo" . "/branches/" . "$release";
my $URL = "$base" . "$repo\n";
my @to = ('your_mail_id@domain.com');
my $LOG = `svn --username=$user  --password=$encode log -r{$dt}:HEAD $URL`;

open LOG,'-|',"svn log -r{$dt}:HEAD $URL" or die $@;
my $i = 0;
while (&ltlog&gt) {
        next if /^----/;
        next if /^$/;
if (/^r/) {
  my($rev, $user) = split /\|/, $_;
                $hash{$rev} = '';
                $users{$rev} = $user;
                } else {
                my @keys = (keys %hash);
                my $key = $keys[$i];
                delete $hash{$key};
                $nuhash{$key} .= $_;
                }
        }

close(LOG);

foreach my $key (keys %hash) {
        if ($hash{$key} =~ /^$/){
        $nuhash{$key} .= '';
          }
        }

foreach my $tab (keys %nuhash) {
                if ($nuhash{$tab} =~ /^$/ || $nuhash{$tab} =~ /^\s+$/ ) {
                $details = `svn log -r$tab $URL`;
                &_send_mail('your_mail_id@domain.com',"$tab :" . "$users{$tab} \
                  => " . " NULL" , @to);
                        }
                }
$i++;
        }
}
#
# Check hash
#
#foreach my $k (keys %nuhash) {
#print "$k" . " => " . "$nuhash{$k}\n";
#}
#
# Send mail to the user
#
sub _send_mail {
my ($from, $sub, @to) = @_;

  my $smtp = Net::SMTP->new('YOUR.SMTP.SERVER');

  $smtp->mail($from);
  $smtp->to(@to);

  $smtp->data();
  $smtp->datasend("To: @to\n");
  $smtp->datasend("Subject: $sub\n");
  $smtp->datasend("\n");
  $smtp->datasend("$details\n");
  $smtp->dataend();

  $smtp->quit;
}

Saturday, January 4, 2014

Secure Tomcat manager for production use

Tomcat manager is very useful for production environment when multiple applications are deployed in a single server. It helps to manage applications without restarting the server.  However, accessing HTML interface of manager application remotely is not a wise decision.
Therefore preventing remote access of tomcat manager using web browser and allowing access of tool-friendly plain text interface instead would be the best choice. This article illustrates a simple solution that has been designed to secure tomcat server for production use.

Tomcat provides a number of Filters to secure the server itself or an individual application. Please check here for more details. Our goal is to prevent web browser to access the Manager application from outside of local host. At the same time we must allow commands as a part of the request URI to get responses in the form of simple text that can be easily parsed and processed. Therefore filter should have logic to allow access based on HTTP request header. A very simple logic could be filtering Remote Address and embedded request properties available in HTTP request header as below.

private String checkHeader = "MyComp";
.
.
.
if (headerValue != null) {
   /*
    * Either connect from 127.0.0.1 or use "tomcatmanager" command
    */
   if (headerValue.equals(checkHeader) || remoteIp.equals("127.0.0.1")) {
    denyStatus = true;
   }
  }

Second part of this solution is a java utility which performs two basic functions. First it encrypts plain text password available in properties file and then decrypt the same again to connect tool-friendly text URI.  Properties file contain plain text user and password as per tomcat-user.xml. Whenever tomcat credential gets change, properties file should get modified accordingly. Another function is to setRequestProperty to prepare URLConnection.

urlConnection.setRequestProperty("referer", "MyComp");

Users with the manager-gui role should not be granted the manager-script or manager-jmx roles. Therefore, to use this client utility, configure tomcat-users.xml accordingly.

Demonstration:
Consider two systems A and B. System A is your Tomcat server where manager application is deployed and system B is your Desktop client. If you try to access HTML interface of tomcat manager from your desktop, it will redirect you to an error page, however if you run the utility it will show you all details in readable text format as below.
  
How it works?
As a client, when you hit web browser to access GUI interface of tomcat manager, filter checks Remote Address and redirect your request. However the filter will allow access of GUI interface from system A as, in such a case, request goes from localhost.

When we run the utility from system B, filter checks and found hardcoded request properties, therefore filter refrain Remote Address checking and allow access of plain text URI.

Saturday, December 14, 2013

X and O puzzle for kids

A very simple game written in python is attached herewith for kids. You can  Download the game and unzip to run on your system.This is perfect for windows platform as a few Windows system commands like “color”, “TMP” file path etc have been used inside the script. However, modifying a few lines the script can be used on other OS like Linux. Please comments if you need source code. The script is converted into exe and zipped to attaché in blog.

Once you double click on the file it will show you the board. Please read the instruction carefully and place your position. Your choice will be placed as "X" and system will place "O" against your choice. 
If none of you win the game the result appears as below.
Different color at end of the game indicates the result. By any chance if you win the game Green Board will congratulate you.
let's enjoy the game and put your comments how is that!!!!


Saturday, August 31, 2013

Search and test all links available in web page

Another small utility to check if all HYPERlinks inside a web page are working operationally. User can modify this utility and use according to their requirement.The nice part of this code is it stores search output in a file and display storage path on GUI console at the end of the process. To test the script simply store the code in a *.py file and double click on it. A popup will appear on desktop. Provide URL/web page and click on "show" button.
Searching will start and result will populate on windows command line console provided, python-3.3.2 is installed and placed in path properly.
#!python
'''
Created on Aug 19, 2013
@requires: Tested in Windows XP
@author: Jaydeb Chakraborty
@version: Python Version-3.3.2

'''
from tkinter import Tk, Label, Button, Entry
import urllib.request
import re
import os
import logging

""" 
Store all URLs in a file and get HTTP return code 
"""
def show():
    strn = entry.get()  
    if re.match('(?:ftp|https)://', strn):
        mesg="Currently HTTPS|FTP is not supported"
        t = Label(w, text=mesg)
        t.pack()
    else:
        l=strn.replace('http://', '')
        mesg='Please check output in ' + (os.environ.get('TEMP', '')) + '\INFO.log'
        t = Label(w, text=mesg)
        t.pack()                
         
    if strn:  
        stdinfo=((os.environ.get('TEMP', ''))+ '\INFO.log')
        logging.basicConfig(filename='%s' % stdinfo, format='%(asctime)s %(message)s', level=logging.INFO)
        logging.warning('*******  Accessing : %s' % strn + ' ********')
        logging.warning('**********************************************')
        local_filename, headers = urllib.request.urlretrieve('http://' + l)
        f = open(local_filename) 
        for lines in f:
            myString_list = [item for item in lines.split(" ")]
            for item in myString_list:      
                try:
                    o = re.search("(?Phttp?://[^\s]+)", item.expandtabs()).group("url")
                    url = re.sub(r'\?.*|".*', "", o)
                    conn = urllib.request.urlopen(url)
                    access = conn.getcode()
                    logging.warning('URL : %s' % url + ' -- Returncode is : %s' % access)
                    print(url, access)
                except :
                    pass 
                
       
w = Tk()
quitBotton = Button(w, text='Quit', command=quit).pack()
showBotton = Button(w, text='Show', command=show).pack()
Label(w, text="        Please provide URL...        ").pack()
entry = Entry(w)
entry.pack()
res = Label(w)
res.pack()
w.title('Test Links in web page')
w.maxsize(1000, 40000)
w.mainloop()

Saturday, August 24, 2013

Windows schedule task: Restart Tomcat automatically

Managing tomcat with a bunch of application sometime becomes problematic. Without investing proper time for root cause analysis it is not possible to isolate the issue perhaps the application from the container which is troubling. Be it a production environment project cannot move forward w farther without resolving the issue. Well, I am talking about development environment where usually we do not think about the impact instead just bounce the server.
For UNIX/Linux it is easy at least can be managed using a few line of script and scheduling the job via cron. Maybe windows expert will do the same very easily. I thought to do the same thing via python. There are many reason and corresponding indication of hanging a server. Having a close look in server log one can found the common foot print that a server left every time it goes to hang state. High Memory utilization is one of them and very common.  Following steps demonstrate how to setup a job written in python in windows.

How to schedule jobs in windows?
>>schtasks /create /tn "PythonCron Job" /tr "path/to/your/executable Arg1 Arg2" /sc hourly
Above command will create job in "Control Panel--> Schedule Task".
In windows, schedule python script does not accept argument.
Therefore converted .py script to .exe file

How to convert .py to .exe?
To Convert .py to .exe for Python3.3 I used CX-Freeze.
For more details please have a look in http://cx-freeze.sourceforge.net/index.html
>>cxfreeze hello.py --target-dir dist
.........and schedule task as
>>C:\Path\To\Your\File.exe Arg1 Arg2

How to adjust schedule job in windows?
Open the task and adjust time in Advance Tab. Every task requires a user name and password
Administrator can use "NT AUTHORITY\SYSTEM"
How to check log?
To check if schedule job is running perfectly, go to "Control Panel--> Schedule Task"
Logs are usually available in "C:\WINDOWS\SchedLgU"
The script also populates logs ‘stdout/stderror’ in user "HOMEPATH". You can check logs if there is any error.
How to check login user in windows
>>echo %USERDOMAIN%\%USERNAME%

How to Test script in DEV?
To test this script please follows the following steps:
    1) Defile CATALINA_HOME in User Variable, if not defined.
    2) Start tomcat server
    3) Execute script in command line "python /path/to/the/script.py MAXMEM TOMCAT_PORT"
    4) For testing provide MAXMEM as low as for example '5'
        Example: >>>python /path/to/the/script.py 5 8080

#!python
'''
Python Version-3.3.2
Created on Aug 19, 2013
@author: Jaydeb Chakraborty
??? Restart Tomcat automatically
??? Monitor server behavior and modify script for different threshold 

!!! Assume CATALINA_HOME is defined and Tomcat starts via 'startup.bat'
'''
import logging
import sys
import re
import os
import socket
import subprocess
import urllib.request
from datetime import datetime
from time import sleep

class ServerProcess:
    def __init__(self, portnum):
        self.portnum = portnum
        self.port = (str(portnum).replace("'", "").replace("]", ""))
    def getPID(self):
        try:
            port = self.port
            p1 = subprocess.Popen('netstat -ano', stdout=subprocess.PIPE)
            p2 = subprocess.Popen('findstr :%s' % port, stdin=p1.stdout, stdout=subprocess.PIPE)
            result=p2.communicate()[0].split()
            pid=str(result[4]).lstrip("b,").replace(",", "").replace("'", "")
            return pid 
        except LookupError as l:
            logging.basicConfig(filename='%s' % stderror, format='%(asctime)s %(message)s', level=logging.DEBUG)
            logging.warning('Watch out! : %s' % l)
            return 0       
    def getMemory(self, pid):
            p3 = subprocess.Popen('tasklist /fi "MEMUSAGE ge 10" /fi "PID eq %s' % pid, stdout=subprocess.PIPE)
            mem = str(p3.communicate()[0].split()[17]).lstrip("b,").replace(",", "").replace("'", "")
            return mem
    def getConnect(self):
        try:
            port = int(self.port)
            sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
            remoteServerIP  = socket.gethostbyname(socket.gethostname())
            conn = sock.connect((remoteServerIP, port))
            sock.close()
            return 'Listening'
        except socket.gaierror as name:
            logging.basicConfig(filename='%s' % stderror, format='%(asctime)s %(message)s', level=logging.DEBUG)
            logging.warning('Watch out! : %s' % name)
            sys.exit()
        except socket.error as s:
            logging.basicConfig(filename='%s' % stderror, format='%(asctime)s %(message)s', level=logging.DEBUG)
            logging.warning('Watch out! : %s' % s)
            return 'NotListening'
            sys.exit()
    def getHttpCode (self):
        port = self.port
        req=('http://localhost:'+ (str(Args[2]).replace(",", "").replace("'", "").replace("]", "")))
        status = urllib.request.urlopen(req)
        data=status.getcode()
        return data        
    def setLogentries(loglevel):
        if (loglevel == 'DEBUG'):
            logging.basicConfig(filename='%s' % stderror, format='%(asctime)s %(message)s', level=logging.DEBUG)
        else:
            logging.basicConfig(filename='%s' % stdout, format='%(asctime)s %(message)s', level=logging.INFO)

if __name__ == "__main__":
    global stderror, stdout 
    stdout=((os.environ.get('HOMEPATH', ''))+ '\INFO.log')
    stderror=((os.environ.get('HOMEPATH', ''))+ '\ERROR.log')
    Args = str(sys.argv).split()
    # Check Arguments
    if ((len(Args)) != 3):
        ServerProcess.setLogentries('DEBUG')
        logging.warning('Watch out! : Provide Hostname and Port name in arguments')
        print("Uses : python " + str(Args[0]).replace("['", "") + " MAXMEM PORT")
        sys.exit()

maxmem=str(Args[1]).replace(",", "").replace("'", "")        
process = ServerProcess(Args[2])
pid = process.getPID()
memory = process.getMemory(pid)      

for i in range(3):
    sleep(1)
    sys.stdout.flush()
    conn = process.getConnect()
    data = process.getHttpCode()
    """ If socket is listening and used memory is less than MAXMEM """
    if(( conn == 'Listening' ) and eval(maxmem)>eval(memory)  and ( data == 200 )):
        ServerProcess.setLogentries('INFO')
        logging.warning('Watch out! : Socket is %s ' % conn + 'when used memory is %s ' % memory + 'and return code is %s ' % data)
    elif(eval(maxmem) > eval(memory) and data == 200):
        ServerProcess.setLogentries('DEBUG')
        logging.warning('Watch out! : Can access url')              
    else:
        p = subprocess.Popen('tasklist /fi "MEMUSAGE ge 10" /fi "PID eq %s' % pid, stdout=subprocess.PIPE)
        proc = str(p.communicate()[0].split()[13]).lstrip("b,").replace(",", "").replace("'", "")
        os.system('taskkill /f /im %s' % proc)
        try:
            startup=(os.environ.get('CATALINA_HOME', '')+'\/bin\/startup.bat')
            os.system('%s' % startup)
            sys.exit()
        except OSError as sys:
            ServerProcess.setLogentries('DEBUG')
            logging.warning('Watch out! : Socket is %s' % sys)